MCPcopy Create free account

hub / github.com/IBM/sarama / functions

Functions3,668 in github.com/IBM/sarama

↓ 1 callersMethodnewHighWatermark
(hwm int)
async_producer.go:914
↓ 1 callersFunctionnewMockMessage
(key, msg Encoder)
mockresponses.go:285
↓ 1 callersMethodnewPartitionOffsetManager
(topic string, partition int32)
offset_manager.go:658
↓ 1 callersMethodnewPartitionProducer
(topic string, partition int32)
async_producer.go:789
↓ 1 callersFunctionnewProduceResponsePromise
()
broker_test.go:94
↓ 1 callersFunctionnewProducerProvider
(brokers []string, producerConfigurationProvider func() *sarama.Config)
examples/txn_producer/main.go:174
↓ 1 callersFunctionnewProducerProvider
(brokers []string, producerConfigurationProvider func() *sarama.Config)
examples/exactly_once/main.go:320
↓ 1 callersFunctionnewReadReplicaTest
(t *testing.T, testConfig readReplicaTestConfig)
consumer_test.go:987
↓ 1 callersFunctionnewSingleFlightRefresher
(f func(topics []string) error)
metadata.go:217
↓ 1 callersFunctionnewSubscriptionMetadataTestGroup
(userData []byte, version KafkaVersion)
consumer_group_test.go:409
↓ 1 callersFunctionnewTestStatefulStrategy
(t *testing.T)
functional_consumer_group_test.go:650
↓ 1 callersMethodnewTopicProducer
(topic string)
async_producer.go:685
↓ 1 callersFunctionnewZstdEncoder
(params ZstdEncoderParams)
zstd.go:50
↓ 1 callersMethodnextUnmuteSignal
nextUnmuteSignal returns the channel that the next unmute (or close) will close.
async_producer.go:230
↓ 1 callersMethodnotifyError
notifyError delivers an abort error and queues a redispatch unless shutdown is already in progress
consumer.go:502
↓ 1 callersMethodoffset
Provide the current offset to record the batch size metric
packet_encoder.go:40
↓ 1 callersFunctionownedPartitions
ownedPartitions sorts the wire representation to keep subscriptions stable
consumer_group_members.go:110
↓ 1 callersFunctionparseCompression
(scheme string)
tools/kafka-producer-performance/main.go:160
↓ 1 callersMethodparseMessages
(msgSet *MessageSet)
consumer.go:809
↓ 1 callersFunctionparsePartitioner
(scheme string, partition int)
tools/kafka-producer-performance/main.go:176
↓ 1 callersMethodparseRecords
(batch *RecordBatch)
consumer.go:843
↓ 1 callersFunctionparseVersion
(version string)
tools/kafka-producer-performance/main.go:195
↓ 1 callersMethodpartialSize
partialSize reports the total on-wire size required to fetch the partial trailing batch fully. Returns 0 when not partial or when the size is unknown
records.go:189
↓ 1 callersFunctionpartitionAndAssert
(t *testing.T, partitioner Partitioner, numPartitions int32, testCase partitionerTestCase)
partitioner_test.go:38
↓ 1 callersFunctionpartitionIDs
(assignment [][]int32)
examples/alter_partition_reassignments/main.go:231
↓ 1 callersMethodpartitionMessage
(msg *ProducerMessage)
async_producer.go:722
↓ 1 callersMethodperformReassignments
Reassign all topic partitions that need reassignment until balanced.
balance_strategy.go:554
↓ 1 callersFunctionplanPartitions
(plan BalanceStrategyPlan, memberID, topic string)
balance_strategy_cooperative_sticky_test.go:160
↓ 1 callersMethodpreferredReadReplicaLease
()
consumer.go:657
↓ 1 callersFunctionprepareDockerTestEnvironment
(ctx context.Context, env *testEnvironment)
functional_test.go:143
↓ 1 callersFunctionprepareTestTopics
(ctx context.Context, env *testEnvironment)
functional_test.go:335
↓ 1 callersFunctionprintErrorAndExit
(code int, format string, values ...any)
tools/kafka-console-partitionconsumer/kafka-console-partitionconsumer.go:91
↓ 1 callersFunctionprintErrorAndExit
(code int, format string, values ...any)
tools/kafka-console-producer/kafka-console-producer.go:143
↓ 1 callersFunctionprintErrorAndExit
(code int, format string, values ...any)
tools/kafka-console-consumer/kafka-console-consumer.go:152
↓ 1 callersFunctionprintUsageErrorAndExit
(format string, values ...any)
tools/kafka-console-partitionconsumer/kafka-console-partitionconsumer.go:97
↓ 1 callersFunctionprintUsageErrorAndExit
(message string)
tools/kafka-console-producer/kafka-console-producer.go:149
↓ 1 callersFunctionprintUsageErrorAndExit
(format string, values ...any)
tools/kafka-console-consumer/kafka-console-consumer.go:158
↓ 1 callersMethodprocessPartitionMovement
Track the movement of a topic partition after assignment
balance_strategy.go:623
↓ 1 callersFunctionprodMsg2Str
(prodMsg *ProducerMessage)
functional_consumer_test.go:493
↓ 1 callersFunctionproduceKeyedWithJava
produceKeyedWithJava produces key:value messages via the Java console producer, which uses the DefaultPartitioner (murmur2) when a key is present.
functional_java_interop_test.go:309
↓ 1 callersFunctionproduceTestRecord
(producerProvider *producerProvider)
examples/txn_producer/main.go:115
↓ 1 callersFunctionproduceWithJava
(t *testing.T, topic string, codec CompressionCodec, messages []string)
functional_java_interop_test.go:35
↓ 1 callersFunctionproduceWithSarama
(t *testing.T, topic string, codec CompressionCodec, messages []string)
functional_java_interop_test.go:114
↓ 1 callersMethodputInt64
(in int64)
real_encoder.go:40
↓ 1 callersMethodputInt8
primitives
real_encoder.go:25
↓ 1 callersMethodputNullableInt32Array
(in []int32)
packet_encoder.go:36
↓ 1 callersMethodputVarint
(in int64)
packet_encoder.go:18
↓ 1 callersMethodputVarint
(in int64)
prep_encoder.go:40
↓ 1 callersMethodputVarint
(in int64)
real_encoder.go:45
↓ 1 callersMethodreadPackage
readPackage reads payload length (4 bytes) and then reads the payload into []byte
gssapi_kerberos.go:82
↓ 1 callersMethodreadToBytes
(r io.Reader)
mockbroker.go:200
↓ 1 callersMethodreassignPartitionToNewConsumer
Identify a new consumer for a topic partition and reassign it.
balance_strategy.go:605
↓ 1 callersMethodregisterCounter
(name string)
broker.go:2029
↓ 1 callersMethodregisterMetrics
()
broker.go:2005
↓ 1 callersMethodrelease
()
offset_manager.go:755
↓ 1 callersMethodrelease
(withCleanup bool)
consumer_group_session.go:272
↓ 1 callersMethodrelease
(producer sarama.AsyncProducer)
examples/txn_producer/main.go:212
↓ 1 callersMethodrelease
(topic string, partition int32, producer sarama.AsyncProducer)
examples/exactly_once/main.go:359
↓ 1 callersFunctionreleaseEncoder
(params ZstdEncoderParams, enc *zstd.Encoder)
zstd.go:70
↓ 1 callersFunctionreleaseLengthField
(m *lengthField)
length_field.go:24
↓ 1 callersMethodremoveMovementRecordOfPartition
(partition topicPartitionAssignment)
balance_strategy.go:995
↓ 1 callersMethodremovePartitions
removePartitions closes partitions revoked during a rebalance, committing first when configured lint:ignore U1000 // consumed by the cooperative reba
offset_manager.go:501
↓ 1 callersMethodrequiredVersion
()
metadata_request.go:185
↓ 1 callersMethodresolveCanonicalNames
(addrs []string)
client.go:1278
↓ 1 callersMethodresponseFeeder
()
consumer.go:747
↓ 1 callersMethodreturnSuccesses
(batch []*ProducerMessage)
async_producer.go:1728
↓ 1 callersMethodrun
()
async_producer.go:1119
↓ 1 callersFunctionrunAsyncProducer
(topic string, partition, messageLoad, messageSize int, config *sarama.Config, brokers []string, throughput i
tools/kafka-producer-performance/main.go:318
↓ 1 callersFunctionrunSyncProducer
(topic string, partition, messageLoad, messageSize, routines int, config *sarama.Config, brokers []string, th
tools/kafka-producer-performance/main.go:365
↓ 1 callersFunctionscanKafkaVersion
(s string, pattern *regexp.Regexp, format string, v [3]*uint)
utils.go:344
↓ 1 callersMethodsendAndReceiveApiVersions
(v int16)
broker.go:1423
↓ 1 callersMethodsendAndReceiveKerberos
()
broker.go:1566
↓ 1 callersMethodsendAndReceiveSASLOAuth
sendAndReceiveSASLOAuth performs the authentication flow as described by KIP-255 https://cwiki.apache.org/confluence/pages/viewpage.action?pageId=7596
broker.go:1707
↓ 1 callersMethodsendAndReceiveSASLPlainAuthV0
In SASL Plain, Kafka expects the auth header to be in the following format Message format (from https://tools.ietf.org/html/rfc4616): message = [a
broker.go:1651
↓ 1 callersMethodsendAndReceiveSASLPlainAuthV1
Kafka 1.x.x onward added a SaslAuthenticate request/response message which wraps the SASL flow in the Kafka protocol, which allows for returning meani
broker.go:1695
↓ 1 callersMethodsendAndReceiveSASLSCRAMv0
()
broker.go:1731
↓ 1 callersMethodsendAndReceiveSASLSCRAMv1
(authSendReceiver func(authBytes []byte) (*SaslAuthenticateResponse, error), scramClient SCRAMClient)
broker.go:1792
↓ 1 callersFunctionsendOffsetCommit
(coordinator *Broker, req *OffsetCommitRequest)
offset_manager.go:323
↓ 1 callersMethodserverError
(err error)
mockbroker.go:395
↓ 1 callersMethodserverLoop
()
mockbroker.go:170
↓ 1 callersFunctionsessionCauseToReason
(cause error)
consumer_group_session.go:381
↓ 1 callersMethodsetGeneration
lint:ignore U1000 // consumed by the cooperative rebalancing path added in a following PR; the only in-build caller here is the unit test (excluded by
offset_manager.go:262
↓ 1 callersMethodsetPartitionCache
(topic string, partitionSet partitionType)
client.go:854
↓ 1 callersMethodsetThrottle
(throttleTime time.Duration)
broker.go:1976
↓ 1 callersMethodsetTypeFromMagic
(pd packetDecoder)
records.go:68
↓ 1 callersMethodsetVersion
(v int16)
offset_fetch_request.go:25
↓ 1 callersFunctionshouldIgnoreMsg
(msg *sarama.ProducerMessage)
examples/interceptors/trace_interceptor.go:43
↓ 1 callersMethodshutdown
()
async_producer.go:1253
↓ 1 callersFunctionsocketErrorProbeAvailable
()
client_test.go:24
↓ 1 callersMethodstart
start starts a new refresh. The refresh is started in a new goroutine, and this function returns a channel on which the caller can wait for the refres
metadata.go:164
↓ 1 callersFunctionstdinAvailable
()
tools/kafka-console-producer/kafka-console-producer.go:157
↓ 1 callersMethodsubscriptionVersion
subscriptionVersion is based on the broker version rather than the rebalance protocol so owned partitions are available during a rolling upgrade
consumer_group.go:591
↓ 1 callersFunctionsupportedProtocols
(strategy BalanceStrategy)
rebalance_protocol.go:43
↓ 1 callersMethodsyncGroupRequest
( coordinator *Broker, members map[string]ConsumerGroupMemberMetadata, plan BalanceStrategyPlan, generatio
consumer_group.go:646
↓ 1 callersMethodtakePartitions
(predicate func(topic string, partition int32) bool)
produce_set.go:134
↓ 1 callersFunctiontestConsumerInterceptor
( t *testing.T, interceptors []ConsumerInterceptor, expectationFn func(*testing.T, int, *ConsumerMessage),
consumer_test.go:2207
↓ 1 callersFunctiontestDecode
(t *testing.T, tp string, key []byte, value []byte)
control_record_test.go:34
↓ 1 callersFunctiontestFuncConsumerGroupFuzzySeed
(topic string)
functional_consumer_group_test.go:316
↓ 1 callersFunctiontestFuncConsumerGroupProduceMessage
--------------------------------------------------------------------
functional_consumer_staticmembership_test.go:219
↓ 1 callersFunctiontestMain
(m *testing.M)
functional_test.go:69
← previousnext →1,201–1,300 of 3,668, ranked by callers