Code
Hub
Workspaces
Following
Trending
Connect
MCP
copy
Create free account
hub
/
github.com/IBM/sarama
/ functions
Functions
3,668 in github.com/IBM/sarama
⨍
Functions
3,668
◇
Types & classes
511
↓ 2 callers
Method
publishTxnPartitions
Makes a request to kafka to add a list of partitions to the current transaction.
transaction_manager.go:785
↓ 2 callers
Method
putArrayLength
(in int)
prep_encoder.go:199
↓ 2 callers
Method
putArrayLength
(in int)
real_encoder.go:204
↓ 2 callers
Method
putFloat64
(in float64)
packet_encoder.go:20
↓ 2 callers
Method
putInt64Array
(in []int64)
packet_encoder.go:35
↓ 2 callers
Method
putRawBytes
(in []byte)
prep_encoder.go:93
↓ 2 callers
Method
putString
(in string)
prep_encoder.go:114
↓ 2 callers
Method
putString
(in string)
prep_encoder.go:213
↓ 2 callers
Method
putString
(in string)
real_encoder.go:109
↓ 2 callers
Method
putString
(in string)
real_encoder.go:219
↓ 2 callers
Method
putVarintBytes
(in []byte)
packet_encoder.go:28
↓ 2 callers
Method
randomizeSeedBrokers
private broker management helpers
client.go:680
↓ 2 callers
Method
reassignPartition
Reassign a specific partition to a new consumer
balance_strategy.go:615
↓ 2 callers
Function
recordBatchTestCases
()
record_test.go:15
↓ 2 callers
Method
registerBroker
registerBroker makes sure a broker received by a Metadata or Coordinator request is registered in the brokers map. It returns the broker that is regis
client.go:747
↓ 2 callers
Method
registerForAllBrokers
(broker *Broker, validator *metricValidator)
metrics_helpers_test.go:29
↓ 2 callers
Method
registerForGlobalAndTopic
(topic string, validator *metricValidator)
metrics_helpers_functional_test.go:11
↓ 2 callers
Method
registerHistogram
(name string)
broker.go:2024
↓ 2 callers
Method
registerMeter
(name string)
broker.go:2019
↓ 2 callers
Function
releaseCrc32Field
(c *crc32Field)
crc32_field.go:29
↓ 2 callers
Method
releasePOMs
Releases/removes closed POMs once they are clean (or when forced)
offset_manager.go:535
↓ 2 callers
Method
releaseSelectedPOMs
releaseSelectedPOMs holds the write lock while closing POM error channels
offset_manager.go:540
↓ 2 callers
Method
removeChild
(child *partitionConsumer)
consumer.go:250
↓ 2 callers
Method
requests
()
offset_manager_test.go:950
↓ 2 callers
Method
reserveLength
()
length_field.go:83
↓ 2 callers
Method
resize
resize rebuilds buf with head back at index 0, sized to 2*count but never below minLen so the zero value of Queue grows on first Add.
internal/queue/queue.go:22
↓ 2 callers
Method
resurrectDeadBrokers
()
client.go:781
↓ 2 callers
Method
retryJoinSync
(ctx context.Context, topics []string, held *heldAssignment, retries int, refreshCoordinator bool)
consumer_group.go:281
↓ 2 callers
Method
retryMessage
(msg *ProducerMessage, err error)
async_producer.go:1738
↓ 2 callers
Method
retryMessages
(batch []*ProducerMessage, err error)
async_producer.go:1747
↓ 2 callers
Function
runConsumerLeaderRefreshErrorTestWithConfig
(t *testing.T, config *Config)
consumer_test.go:416
↓ 2 callers
Method
safelyApplyInterceptor
(interceptor ProducerInterceptor)
interceptors.go:25
↓ 2 callers
Method
satisfies
(other kafkaVersion)
functional_test.go:556
↓ 2 callers
Method
send
b.lock must be held by caller a non-nil res results in a response promise being created
broker.go:1119
↓ 2 callers
Method
sendAndReceiveSASLHandshake
(saslType SASLMechanism, version int16)
broker.go:1574
↓ 2 callers
Method
sendInternal
b.lock must be held by caller
broker.go:1163
↓ 2 callers
Method
sendWithPromise
b.lock must be held by caller
broker.go:1144
↓ 2 callers
Function
setupToxiProxies
setupToxiProxies will configure the toxiproxy proxies with routes for the kafka brokers if they don't already exist
functional_test.go:120
↓ 2 callers
Function
sortPartitions
(currentAssignment map[string][]topicPartitionAssignment, partitionsWithADifferentPreviousAssignment map[topic
balance_strategy.go:714
↓ 2 callers
Function
sortPartitionsByPotentialConsumerAssignments
(partition2AllPotentialConsumers map[topicPartitionAssignment][]string)
balance_strategy.go:794
↓ 2 callers
Method
spn
(broker *Broker)
gssapi_kerberos.go:201
↓ 2 callers
Method
stopConsuming
()
consumer.go:1386
↓ 2 callers
Method
subscriptionMetadata
subscriptionMetadata builds the ConsumerGroupMemberMetadata for a single strategy in a JoinGroup request. If the strategy implements SubscriptionUserD
consumer_group.go:607
↓ 2 callers
Function
tearDownDockerTestEnvironment
(ctx context.Context, env *testEnvironment)
functional_test.go:284
↓ 2 callers
Function
testRequestWithoutByteComparison
(t *testing.T, name string, rb protocolBody)
request_test.go:605
↓ 2 callers
Function
topicWithEvenLeaders
(t *testing.T, adminClient ClusterAdmin, client Client, numPartitions int32)
functional_admin_test.go:20
↓ 2 callers
Method
updateBroker
(brokers []*Broker)
client.go:713
↓ 2 callers
Method
updateLeader
()
async_producer.go:963
↓ 2 callers
Method
updateRequestLatencyAndInFlightMetrics
(requestLatency time.Duration)
broker.go:1903
↓ 2 callers
Function
validServerNameTLS
(addr string, cfg *tls.Config)
broker.go:2034
↓ 2 callers
Function
validateGroupInstanceId
(id string)
config.go:946
↓ 2 callers
Function
verifyProducerConfig
(config *Config)
sync_producer.go:118
↓ 2 callers
Function
verifyValidityAndBalance
(t *testing.T, consumers map[string]ConsumerGroupMemberMetadata, plan BalanceStrategyPlan)
balance_strategy_test.go:2187
↓ 2 callers
Function
version
()
version.go:14
↓ 2 callers
Method
wait
wait returns the channel on which you can wait for the refresh to complete. You need to hold the lock to call this method.
metadata.go:200
↓ 2 callers
Method
waitIfThrottled
()
broker.go:1988
↓ 2 callers
Method
waitUntilMuted
waitUntilMuted blocks until all partitions in the set can be muted, then mutes them. Returns false if the muter was closed before all partitions could
async_producer.go:206
↓ 2 callers
Function
zstdCompress
(params ZstdEncoderParams, dst, src []byte)
zstd.go:119
↓ 1 callers
Method
AddBlock
(topic string, partition int32, block *OffsetFetchResponseBlock)
offset_fetch_response.go:68
↓ 1 callers
Method
AddBlock
(topic string, partitionID int32, timestamp int64, maxOffsets int32)
offset_request.go:258
↓ 1 callers
Method
AddControlRecord
(topic string, partition int32, offset int64, producerID int64, recordType ControlRecordType)
fetch_response.go:732
↓ 1 callers
Method
AddControlRecordWithTimestamp
(topic string, partition int32, offset int64, producerID int64, recordType ControlRecordType, timestamp time.T
fetch_response.go:686
↓ 1 callers
Method
AddError
(topic string, partition int32, kerror KError, message *string)
alter_partition_reassignments_response.go:46
↓ 1 callers
Method
AddError
(topic string, partition int32, errorCode KError)
delete_offsets_response.go:20
↓ 1 callers
Method
AddGroupAssignmentMember
( memberId string, memberAssignment *ConsumerGroupMemberAssignment, )
sync_group_request.go:199
↓ 1 callers
Method
AddGroupBlock
AddGroupBlock adds a block for groupID/topic/partition, creating the group entry if needed. On v0-7 it falls back to AddBlock and ignores groupID.
offset_fetch_response.go:455
↓ 1 callers
Method
AddGroupDescription
(groupID string, description *GroupDescription)
mockresponses.go:110
↓ 1 callers
Method
AddGroupPartition
AddGroupPartition adds a partition under group in a v8+ batch, creating the group entry if absent.
offset_fetch_request.go:361
↓ 1 callers
Method
AddGroupProtocolMetadata
(name string, metadata *ConsumerGroupMemberMetadata)
join_group_request.go:246
↓ 1 callers
Method
AddMessageToTxn
(msg *sarama.ConsumerMessage, groupId string, metadata *string)
mocks/sync_producer.go:267
↓ 1 callers
Method
AddMessageToTxnWithGroupMetadata
(msg *ConsumerMessage, groupMetadata *ConsumerGroupMetadata, metadata *string)
async_producer.go:455
↓ 1 callers
Method
AddOffsetsToTxn
(offsets map[string][]*sarama.PartitionOffsetMetadata, groupId string)
mocks/sync_producer.go:259
↓ 1 callers
Method
AddOffsetsToTxnWithGroupMetadata
AddOffsetsToTxnWithGroupMetadata adds associated offsets to the current transaction, carrying the consumer group member metadata so the broker can fen
sync_producer.go:60
↓ 1 callers
Method
AddPartition
(topic string, partitionID int32)
offset_fetch_request.go:10
↓ 1 callers
Method
AddPartition
(topic string, partitionID int32)
delete_offsets_request.go:95
↓ 1 callers
Method
AddRecordBatch
(topic string, partition int32, key, value Encoder, offset int64, producerID int64, isTransactional bool)
fetch_response.go:728
↓ 1 callers
Method
AddRecordBatchWithTimestamp
AddRecordBatchWithTimestamp is similar to AddRecordWithTimestamp But instead of appending 1 record to a batch, it append a new batch containing 1 reco
fetch_response.go:664
↓ 1 callers
Method
AddSet
(topic string, partition int32, set *MessageSet)
produce_request.go:310
↓ 1 callers
Method
AlterConfigs
AlterConfigs sends a request to alter config and return a response or error
broker.go:932
↓ 1 callers
Method
AssertNoInitialValues
()
functional_consumer_group_test.go:695
↓ 1 callers
Method
AssignmentData
AssignmentData serializes the set of topics currently assigned to the specified member as part of the supplied balance plan
balance_strategy.go:340
↓ 1 callers
Method
AsyncClose
()
offset_manager.go:721
↓ 1 callers
Method
AsyncClose
()
async_producer.go:583
↓ 1 callers
Method
AsyncClose
PartitionConsumer interface implementation AsyncClose implements the AsyncClose method from the sarama.PartitionConsumer interface.
mocks/consumer.go:273
↓ 1 callers
Method
AsyncClose
Implement Producer interface AsyncClose corresponds with the AsyncClose method of sarama's Producer implementation. By closing a mock producer, you
mocks/async_producer.go:128
↓ 1 callers
Method
Authorize
Authorize performs the kerberos auth handshake for authorization
gssapi_kerberos.go:235
↓ 1 callers
Method
AuthorizeV2
AuthorizeV2 performs the SASL v2 GSSAPI authentication with the Kafka broker.
gssapi_kerberos.go:283
↓ 1 callers
Method
BeginTxn
()
mocks/sync_producer.go:231
↓ 1 callers
Method
Claims
Claims returns information about the claimed partitions by topic.
consumer_group_session.go:13
↓ 1 callers
Method
Close
Close stops the OffsetManager from managing offsets. It is required to call this function before an OffsetManager object passes out of scope, as it wi
offset_manager.go:22
↓ 1 callers
Method
Close
()
client.go:312
↓ 1 callers
Method
Close
()
functional_consumer_group_test.go:373
↓ 1 callers
Method
Close
()
examples/http_server/http_server.go:92
↓ 1 callers
Method
Commit
()
offset_manager.go:256
↓ 1 callers
Method
Commit
Commit the offset to the backend Note: calling Commit performs a blocking synchronous operation.
consumer_group_session.go:39
↓ 1 callers
Method
CommitTxn
()
mocks/sync_producer.go:239
↓ 1 callers
Method
Config
()
client.go:252
↓ 1 callers
Method
ConsumeClaim
ConsumeClaim must start a consumer loop of ConsumerGroupClaim's Messages(). Once the Messages() channel is closed, the Handler must finish its process
consumer_group_session.go:423
↓ 1 callers
Method
CreateACLs
Creates access control lists (ACLs) which are bound to specific resources. This operation is not transactional so it may succeed for some ACLs while f
admin.go:109
↓ 1 callers
Method
CreateProxy
( name string, listenAddr string, targetAddr string, )
internal/toxiproxy/client.go:38
← previous
next →
701–800 of 3,668, ranked by callers