MCPcopy Create free account

hub / github.com/apache/pulsar-client-go / functions

Functions5,094 in github.com/apache/pulsar-client-go

↓ 1 callersMethodnewDeviceCodeRequest
newDeviceCodeRequest builds a new DeviceCodeRequest wrapped in an http.Request
oauth2/device_code_provider.go:97
↓ 1 callersFunctionnewKeyBasedBatches
newKeyBasedBatches init a keyBasedBatches
pulsar/internal/key_based_batch_builder.go:59
↓ 1 callersMethodnewLookupService
(serviceURL string)
pulsar/internal/rpc_client.go:243
↓ 1 callersFunctionnewMockMetrics
newMockMetrics creates Metrics with real prometheus counters for testing.
pulsar/internal/connection_test.go:214
↓ 1 callersFunctionnewMultiTopicConsumer
(client *client, options ConsumerOptions, topics []string, messageCh chan ConsumerMessage, dlq *dlqRouter, rl
pulsar/consumer_multitopic.go:84
↓ 1 callersFunctionnewPartitionExpansionRaceConnection
()
pulsar/consumer_test.go:6158
↓ 1 callersFunctionnewPartitionProducer
(client *client, topic string, options *ProducerOptions, partitionIdx int, metrics *internal.LeveledMetrics)
pulsar/producer_partition.go:153
↓ 1 callersFunctionnewProducer
(client *client, options *ProducerOptions)
pulsar/producer_impl.go:72
↓ 1 callersFunctionnewProducerCommand
()
perf/perf-producer.go:46
↓ 1 callersFunctionnewSchemaCache
()
pulsar/producer_partition.go:136
↓ 1 callersFunctionnewSchemaInfoCache
(client *client, topic string)
pulsar/consumer_partition.go:359
↓ 1 callersFunctionnewTableView
(client *client, options TableViewOptions)
pulsar/table_view_impl.go:55
↓ 1 callersFunctionnewTestReconnectBackoffPolicy
(minBackoff, maxBackoff time.Duration)
pulsar/producer_test.go:2963
↓ 1 callersFunctionnewTransactionCoordinatorClientImpl
newTransactionCoordinatorClientImpl init a transactionImpl coordinator client and acquire connections with all transactionImpl coordinators.
pulsar/transaction_coordinator_client.go:338
↓ 1 callersMethodnewTransactionHandler
(partition uint64)
pulsar/transaction_coordinator_client.go:68
↓ 1 callersFunctionnewUnAckChunksTracker
(pc *partitionConsumer)
pulsar/consumer_partition.go:2821
↓ 1 callersFunctionnewUnexpectedErrMsg
NewUnexpectedErrMsg instantiates an ErrUnexpectedMsg error. Optionally provide a list of IDs associated with the message for additional context in the
pulsar/helper.go:34
↓ 1 callersFunctionnewZeroConsumer
(client *client, options ConsumerOptions, topic string, messageCh chan ConsumerMessage, dlq *dlqRouter, rlq
pulsar/consumer_zero_queue.go:55
↓ 1 callersMethodnextTCNumber
()
pulsar/transaction_coordinator_client.go:447
↓ 1 callersFunctionnormalizeValidatedHost
(hostname, host string)
pulsar/internal/service_uri.go:232
↓ 1 callersFunctionparsePackageType
(packageTypeName string)
pulsaradmin/pkg/utils/package_type.go:30
↓ 1 callersFunctionparseParams
(params string)
pulsar/auth/provider.go:91
↓ 1 callersFunctionparseTopicSchemaCompatibilityStrategy
(body []byte)
pulsaradmin/pkg/admin/topic.go:2393
↓ 1 callersMethodpeekNthMessage
( ctx context.Context, topic utils.TopicName, sName string, pos int, )
pulsaradmin/pkg/admin/subscription.go:304
↓ 1 callersMethodperiodicPartitionUpdateCheck
()
pulsar/table_view_impl.go:148
↓ 1 callersMethodprepareTransaction
(sr *sendRequest)
pulsar/producer_partition.go:1164
↓ 1 callersMethodprev
()
pulsar/impl_message.go:119
↓ 1 callersMethodprocessMessageChunk
(compressedPayload internal.Buffer, msgMeta *pb.MessageMetadata, pbMsgID *pb.MessageIdData)
pulsar/consumer_partition.go:1526
↓ 1 callersFunctionproduce
(produceArgs *ProduceArgs, stop <-chan struct{})
perf/perf-producer.go:80
↓ 1 callersMethodreadChecksum
ReadChecksum
pulsar/internal/commands.go:83
↓ 1 callersFunctionreadElement
(r io.Reader, element interface{})
pulsar/primitiveSerDe.go:224
↓ 1 callersMethodreadFromConnection
()
pulsar/internal/connection_reader.go:40
↓ 1 callersMethodreadMessage
()
pulsar/internal/commands.go:158
↓ 1 callersMethodreadSingleMessage
()
pulsar/internal/commands.go:165
↓ 1 callersMethodreceivedCommand
(cmd *pb.BaseCommand, headersAndPayload Buffer)
pulsar/internal/connection.go:553
↓ 1 callersMethodreconnectToBroker
()
pulsar/transaction_coordinator_client.go:148
↓ 1 callersMethodreconnectToBroker
(connectionClosed *connectionClosed)
pulsar/consumer_partition.go:2100
↓ 1 callersMethodrequestSeek
(msgID *messageID)
pulsar/consumer_partition.go:1084
↓ 1 callersMethodreserveMem
(sr *sendRequest)
pulsar/producer_partition.go:1740
↓ 1 callersMethodreserveResources
(sr *sendRequest)
pulsar/producer_partition.go:1762
↓ 1 callersMethodreserveSemaphore
(sr *sendRequest)
pulsar/producer_partition.go:1712
↓ 1 callersMethodreset
()
pulsar/ack_grouping_tracker_test.go:91
↓ 1 callersMethodreset
()
pulsar/internal/batch_builder.go:245
↓ 1 callersMethodreset
()
pulsar/internal/key_based_batch_builder.go:180
↓ 1 callersFunctionrespIsOk
respIsOk is used to validate a successful http status code
pulsar/internal/http_client.go:267
↓ 1 callersFunctionrespIsOk
respIsOk is used to validate a successful http status code
pulsaradmin/pkg/rest/client.go:501
↓ 1 callersFunctionresponseError
responseError is used to parse a response into a client error
pulsar/internal/http_client.go:318
↓ 1 callersFunctionresponseError
responseError is used to parse a response into a client error
pulsaradmin/pkg/rest/client.go:570
↓ 1 callersMethodrun
()
pulsar/retry_router.go:95
↓ 1 callersMethodrun
()
pulsar/dlq_router.go:106
↓ 1 callersMethodrun
()
pulsar/internal/connection.go:394
↓ 1 callersFunctionrunAckIDListTest
(t *testing.T, enableBatchIndexAck bool)
pulsar/consumer_test.go:5611
↓ 1 callersMethodrunBackgroundPartitionDiscovery
(period time.Duration)
pulsar/producer_impl.go:159
↓ 1 callersMethodrunBackgroundPartitionDiscovery
(period time.Duration)
pulsar/consumer_impl.go:321
↓ 1 callersFunctionrunBatchIndexAckTest
(t *testing.T, ackWithResponse bool, cumulative bool, option *AckGroupingOptions)
pulsar/consumer_test.go:4825
↓ 1 callersMethodrunEventsLoop
()
pulsar/transaction_coordinator_client.go:126
↓ 1 callersMethodrunEventsLoop
()
pulsar/producer_partition.go:559
↓ 1 callersMethodrunEventsLoop
()
pulsar/consumer_partition.go:2008
↓ 1 callersFunctionrunMultiTopicAckIDList
(t *testing.T, regex bool)
pulsar/consumer_multitopic_test.go:235
↓ 1 callersMethodrunPingCheck
(pingCheckTicker *time.Ticker)
pulsar/internal/connection.go:440
↓ 1 callersMethodscopedPutWithContext
( ctx context.Context, endpoint string, body interface{}, params map[string]string, )
pulsaradmin/pkg/admin/topic_policies.go:213
↓ 1 callersMethodsendIndividualAckWithTxn
(msgID MessageID, txn *transaction)
pulsar/consumer_partition.go:769
↓ 1 callersMethodsendPing
()
pulsar/internal/connection.go:837
↓ 1 callersFunctionserveProfiling
use `http://addr/debug/pprof` to access the browser use `go tool pprof http://addr/debug/pprof/profile` to get pprof file(cpu info) use `go tool pprof
perf/pulsar-perf-go.go:150
↓ 1 callersMethodsetFirstChunkID
(msgID *messageID)
pulsar/producer_partition.go:1823
↓ 1 callersMethodsetLastChunkID
(msgID *messageID)
pulsar/producer_partition.go:1827
↓ 1 callersMethodsetLastDequeuedMsg
(msgID MessageID)
pulsar/consumer_impl.go:839
↓ 1 callersMethodsetNamespaceIsolationPolicy
(ctx context.Context, cluster, policyName string, namespaceIsolationData utils.NamespaceIsolationData)
pulsaradmin/pkg/admin/ns_isolation_policy.go:108
↓ 1 callersMethodsetPrevBatchAcked
()
pulsar/impl_message.go:462
↓ 1 callersMethodsetRunning
(isRunning bool)
pulsar/internal/memory_limit_controller.go:56
↓ 1 callersMethodsetStateReady
()
pulsar/internal/connection.go:1104
↓ 1 callersMethodshouldSendToDlq
(cm *ConsumerMessage)
pulsar/dlq_router.go:80
↓ 1 callersFunctionsplitHostPortOrDefault
(serviceName string, serviceInfos []string, hostname string)
pulsar/internal/service_uri.go:195
↓ 1 callersFunctionsplitHostURI
(uriStr string)
pulsar/internal/service_uri.go:156
↓ 1 callersMethodsubscribe
(topics []string, dlq *dlqRouter, rlq *retryRouter)
pulsar/consumer_regex.go:412
↓ 1 callersFunctiontestCompression
(b *testing.B, provider Provider)
pulsar/internal/compression/compression_bench_test.go:29
↓ 1 callersFunctiontestDecompression
(b *testing.B, provider Provider)
pulsar/internal/compression/compression_bench_test.go:46
↓ 1 callersFunctiontestEndpoint
(parts ...string)
pulsar/helper_for_test.go:58
↓ 1 callersFunctiontestReaderSeekByIDWithHasNext
(t *testing.T, startMessageID MessageID, startMessageIDInclusive bool)
pulsar/reader_test.go:1142
↓ 1 callersFunctiontestReaderSeekByTimeWithHasNext
(t *testing.T, startMessageID MessageID)
pulsar/reader_test.go:1221
↓ 1 callersFunctiontestTopicMigrate
( t *testing.T, blueAdminURL string, blueClientUrl string, greenAdminURL string, migrationBody string)
pulsar/blue_green_migration_test.go:78
↓ 1 callersFunctiontestTopicUnload
(t *testing.T, adminURL string, clientEndpointFunc func(utils.LookupData) string, unloadEndpointFunc func(ut
pulsar/extensible_load_manager_test.go:104
↓ 1 callersMethodtoHTTP
()
pulsar/internal/http_client.go:243
↓ 1 callersMethodtoHTTP
(ctx context.Context)
pulsaradmin/pkg/rest/client.go:477
↓ 1 callersFunctiontoProtoInitialPosition
(p SubscriptionInitialPosition)
pulsar/consumer_impl.go:895
↓ 1 callersFunctiontoProtoKeySharedMeta
(ksp *KeySharedPolicy)
pulsar/key_shared_policy.go:59
↓ 1 callersFunctiontoProtoProducerAccessMode
(accessMode ProducerAccessMode)
pulsar/producer_partition.go:1831
↓ 1 callersFunctiontoProtoSubType
(st SubscriptionType)
pulsar/consumer_impl.go:880
↓ 1 callersFunctiontopicPath
(topic string)
pulsar/helper_for_test.go:167
↓ 1 callersMethodtrack
()
pulsar/negative_acks_tracker.go:149
↓ 1 callersFunctiontrimLowerBit
(ts int64, precisionBit int64)
pulsar/negative_acks_tracker.go:90
↓ 1 callersMethodtryAddIndividual
(id MessageID)
pulsar/ack_grouping_tracker.go:172
↓ 1 callersMethodtryUpdateCumulative
(id MessageID)
pulsar/ack_grouping_tracker.go:191
↓ 1 callersMethodunsubscribe
(topics []string)
pulsar/consumer_regex.go:430
↓ 1 callersMethodupdateChunkInfo
(sr *sendRequest)
pulsar/producer_partition.go:1269
↓ 1 callersMethodupdateMetaData
(sr *sendRequest)
pulsar/producer_partition.go:1243
↓ 1 callersMethodupdateSchema
(sr *sendRequest)
pulsar/producer_partition.go:1189
↓ 1 callersMethodupdateSingleMessageMetadataSeqID
(smm *pb.SingleMessageMetadata, msg *ProducerMessage)
pulsar/producer_partition.go:762
↓ 1 callersMethodupdateUncompressedPayload
(sr *sendRequest)
pulsar/producer_partition.go:1218
↓ 1 callersMethoduseragent
()
pulsar/internal/http_client.go:229
← previousnext →1,601–1,700 of 5,094, ranked by callers