Code
Hub
Workspaces
Following
Trending
Connect
MCP
copy
Create free account
hub
/
github.com/apache/pulsar-client-go
/ functions
Functions
5,094 in github.com/apache/pulsar-client-go
⨍
Functions
5,094
◇
Types & classes
637
↓ 3 callers
Function
sortMessageIDs
(msgIDs []messageID)
pulsar/negative_acks_tracker_test.go:72
↓ 3 callers
Function
subscriber
(c *client, topics []string, opts ConsumerOptions, ch chan ConsumerMessage, dlq *dlqRouter, rlq *retryRouter)
pulsar/consumer_regex.go:478
↓ 3 callers
Method
success
()
pulsar/consumer_test.go:6249
↓ 3 callers
Function
validateTopicNames
(topics ...string)
pulsar/helper.go:56
↓ 3 callers
Function
verifyLogOutput
(t *testing.T, logOutput, expectedLevel, expectedMessage string, expectedFields ...Fields)
pulsar/log/wrapper_slog_test.go:169
↓ 3 callers
Method
verifyOpen
()
pulsar/transaction_impl.go:223
↓ 3 callers
Method
waitWithContext
waitWithContext Same as wait() call, but the end condition can also be controlled through the context. It blocks until either a broadcast occurs or th
pulsar/internal/channel_cond.go:51
↓ 3 callers
Method
writeData
(buffer internal.Buffer, sequenceID uint64, callbacks []interface{})
pulsar/producer_partition.go:904
↓ 2 callers
Method
ActiveConsumerChanged
(isActive bool)
pulsar/internal/connection.go:97
↓ 2 callers
Method
BeforeConsume
BeforeConsume This is called just before the message is send to Consumer's ConsumerMessage channel.
pulsar/consumer_interceptor.go:22
↓ 2 callers
Method
BeforeSend
BeforeSend This is called before send the message to the brokers. This method is allowed to modify the message.
pulsar/producer_interceptor.go:23
↓ 2 callers
Function
CheckName
(name string)
pulsaradmin/pkg/utils/namespace_name.go:86
↓ 2 callers
Method
Close
()
pulsar/consumer_partition.go:1030
↓ 2 callers
Method
Close
()
pulsar/internal/client_handlers.go:52
↓ 2 callers
Function
ConvertGetAllSchemasResponseToSchemaInfosWithVersion
( tn *TopicName, response GetAllSchemasResponse, )
pulsaradmin/pkg/utils/schema_util.go:98
↓ 2 callers
Function
Crc32cCheckSum
Crc32cCheckSum handles computing the checksum.
pulsar/internal/checksum.go:34
↓ 2 callers
Method
CreateNamespace
CreateNamespace creates a new empty namespace with no policies attached
pulsaradmin/pkg/admin/namespace.go:50
↓ 2 callers
Method
CreateNsWithBundlesDataWithContext
( ctx context.Context, namespace string, bundleData *utils.BundlesData, )
pulsaradmin/pkg/admin/namespace.go:846
↓ 2 callers
Method
CreateSchemaByPayloadWithContext
( ctx context.Context, topic string, schemaPayload utils.PostSchemaPayload, )
pulsaradmin/pkg/admin/schema.go:213
↓ 2 callers
Method
CreateWithPropertiesWithContext
( ctx context.Context, topic utils.TopicName, partitions int, meta map[string]string, )
pulsaradmin/pkg/admin/topic.go:1150
↓ 2 callers
Function
DLQWithProducerOptions
(t *testing.T, prodOpt *ProducerOptions)
pulsar/consumer_test.go:1941
↓ 2 callers
Method
DeleteWithContext
DeleteWithContext deletes an existing cluster
pulsaradmin/pkg/admin/cluster.go:50
↓ 2 callers
Method
Descriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:80
↓ 2 callers
Method
Descriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:142
↓ 2 callers
Method
Descriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:271
↓ 2 callers
Method
Descriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:330
↓ 2 callers
Method
Descriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:448
↓ 2 callers
Method
Descriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:504
↓ 2 callers
Method
Descriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:560
↓ 2 callers
Method
Descriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:673
↓ 2 callers
Method
Descriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:735
↓ 2 callers
Method
Descriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:791
↓ 2 callers
Method
Descriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:847
↓ 2 callers
Method
Descriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:906
↓ 2 callers
Method
Descriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:962
↓ 2 callers
Method
Descriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:1030
↓ 2 callers
Method
Descriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:1086
↓ 2 callers
Method
Descriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:1145
↓ 2 callers
Method
Descriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:1370
↓ 2 callers
Method
Encrypt
([]byte, *pb.MessageMetadata)
pulsar/internal/crypto/encryptor.go:26
↓ 2 callers
Function
ExtractSpanContextFromProducerMessage
(message *pulsar.ProducerMessage)
pulsar/internal/pulsartracing/message_carrier_util.go:40
↓ 2 callers
Method
Float64
(buf []byte)
pulsar/primitiveSerDe.go:101
↓ 2 callers
Method
FlushBatches
FlushBatches all the messages buffered in multiple batches and wait until all messages have been successfully persisted.
pulsar/internal/batch_builder.go:60
↓ 2 callers
Method
GetAckSet
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:1549
↓ 2 callers
Method
GetAssignedBrokerServiceUrl
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:4604
↓ 2 callers
Method
GetAssignedBrokerServiceUrlTls
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:4611
↓ 2 callers
Function
GetAuthProvider
(config *config.Config)
pulsaradmin/pkg/admin/auth/provider.go:57
↓ 2 callers
Method
GetBrokerAddress
(brokerServiceURL string, proxyThroughServiceURL bool)
pulsar/internal/lookup_service.go:76
↓ 2 callers
Method
GetBrokerAddress
(brokerServiceURL string, proxyThroughServiceURL bool)
pulsar/internal/lookup_service.go:137
↓ 2 callers
Method
GetBundleRange
GetBundleRange returns a bundle range of a topic
pulsaradmin/pkg/admin/topic.go:225
↓ 2 callers
Function
GetCompressionProvider
( compressionType pb.CompressionType, level compression.Level, )
pulsar/internal/batch_builder.go:308
↓ 2 callers
Method
GetConnections
GetConnections get all connections in the pool.
pulsar/internal/connection_pool.go:38
↓ 2 callers
Function
GetConnectionsCount
(p *ConnectionPool)
pulsar/internal/helper.go:28
↓ 2 callers
Method
GetEncryptionContext
GetEncryptionContext returns the ecryption context of the message. It will be used by the application to parse the undecrypted message.
pulsar/message.go:140
↓ 2 callers
Method
GetLastMessageId
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:5384
↓ 2 callers
Method
GetMaxTopicsPerNamespace
GetMaxTopicsPerNamespace returns the maxTopicsPerNamespace for a namespace.
pulsaradmin/pkg/admin/namespace.go:364
↓ 2 callers
Method
GetMessageID
GetMessageID returns the message Id by timestamp(ms) of a topic @param topic topicName struct @param timestamp absolute timestamp (in ms)
pulsaradmin/pkg/admin/topic.go:242
↓ 2 callers
Method
GetMessagesByIDWithContext
( ctx context.Context, topic utils.TopicName, ledgerID, entryID int64, )
pulsaradmin/pkg/admin/subscription.go:347
↓ 2 callers
Method
GetNullValue
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:2053
↓ 2 callers
Method
GetOrderingKey
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:2011
↓ 2 callers
Function
GetPackageNameWithComponents
(packageType PackageType, tenant, namespace, name, version string)
pulsaradmin/pkg/utils/package_name.go:42
↓ 2 callers
Method
GetPayloadSize
()
pulsaradmin/pkg/utils/message.go:73
↓ 2 callers
Method
GetProducerReady
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:4901
↓ 2 callers
Method
GetSchemaCompatibilityStrategyAppliedWithContext
( ctx context.Context, topic utils.TopicName, applied bool, )
pulsaradmin/pkg/admin/topic.go:2361
↓ 2 callers
Method
GetState
GetState Get the state of the transaction.
pulsar/transaction.go:66
↓ 2 callers
Method
GetStatsWithOptionWithContext
GetStatsWithOptionWithContext returns the stats for the topic All the rates are computed over a 1-minute window and are relative the last completed 1
pulsaradmin/pkg/admin/topic.go:284
↓ 2 callers
Method
GetSubPermissions
GetSubPermissions returns subscription permissions on a namespace
pulsaradmin/pkg/admin/namespace.go:532
↓ 2 callers
Method
GetSubscriptionExpirationTime
GetSubscriptionExpirationTime gets the subscription expiration time on a namespace. Returns -1 if not set
pulsaradmin/pkg/admin/namespace.go:721
↓ 2 callers
Method
GetTLSCertificate
()
pulsar/auth/tls.go:64
↓ 2 callers
Method
GetTokens
(identifier string)
oauth2/config_tokenprovider.go:23
↓ 2 callers
Method
GetTopicAutoCreation
GetTopicAutoCreation returns the topic auto-creation config for a namespace. Returns nil if the topic auto-creation config is not configured at the na
pulsaradmin/pkg/admin/namespace.go:151
↓ 2 callers
Method
GetVersionByPayloadWithContext
( ctx context.Context, topic string, schemaPayload utils.PostSchemaPayload, )
pulsaradmin/pkg/admin/schema.go:259
↓ 2 callers
Method
HasURL
()
pulsar/producer_partition.go:416
↓ 2 callers
Method
HealthCheckWithTopicVersion
HealthCheckWithTopicVersion runs a health check on the broker
pulsaradmin/pkg/admin/brokers.go:100
↓ 2 callers
Method
HealthCheckWithTopicVersionWithContext
(ctx context.Context, topicVersion utils.TopicVersion)
pulsaradmin/pkg/admin/brokers.go:260
↓ 2 callers
Method
Index
Index returns index from broker entry metadata, or empty if the feature is not enabled in the broker.
pulsar/message.go:144
↓ 2 callers
Function
InjectConsumerMessageSpanContext
(ctx context.Context, message pulsar.ConsumerMessage)
pulsar/internal/pulsartracing/message_carrier_util.go:64
↓ 2 callers
Function
InjectProducerMessageSpanContext
(ctx context.Context, message *pulsar.ProducerMessage)
pulsar/internal/pulsartracing/message_carrier_util.go:28
↓ 2 callers
Method
InvalidateToken
InvalidateToken is called when the token is rejected by the resource server.
oauth2/cache/cache.go:35
↓ 2 callers
Method
IsHTTP
()
pulsar/internal/service_uri.go:85
↓ 2 callers
Method
IsMaxBackoffReached
IsMaxBackoffReached evaluates if the max number of retries is reached
pulsar/backoff/backoff.go:32
↓ 2 callers
Method
IsNullValue
IsNullValue reports whether the message was published as a null-value (tombstone) message, i.e. with MessageMetadata.null_value set. For such messages
pulsar/message.go:97
↓ 2 callers
Method
IsProxied
()
pulsar/internal/connection.go:91
↓ 2 callers
Function
MakeHTTPPath
(apiVersion string, componentPath string)
pulsaradmin/pkg/utils/utils.go:26
↓ 2 callers
Method
MakeRequestWithURLWithContext
( ctx context.Context, method string, urlOpt *url.URL, )
pulsaradmin/pkg/rest/client.go:117
↓ 2 callers
Function
NewAuthDisabled
NewAuthDisabled return a interface of Provider
pulsar/auth/disabled.go:28
↓ 2 callers
Function
NewAuthenticationAthenzWithParams
(params map[string]string)
pulsar/auth/athenz.go:70
↓ 2 callers
Function
NewAuthenticationFromTLSCertSupplier
NewAuthenticationFromTLSCertSupplier Create new Authentication provider with specified TLS certificate supplier
pulsar/client.go:66
↓ 2 callers
Function
NewAuthenticationOAuth2WithFlow
( issuer oauth2.Issuer, flowOptions oauth2.ClientCredentialsFlowOptions)
pulsaradmin/pkg/admin/auth/oauth2.go:60
↓ 2 callers
Function
NewAuthenticationTLS
NewAuthenticationTLS initialize the authentication provider
pulsar/auth/tls.go:41
↓ 2 callers
Function
NewAuthenticationTLS
NewAuthenticationTLS initialize the authentication provider
pulsaradmin/pkg/admin/auth/tls.go:44
↓ 2 callers
Function
NewAuthenticationToken
NewAuthenticationToken returns a token auth provider that will use the specified token to talk with Pulsar brokers
pulsar/auth/token.go:47
↓ 2 callers
Function
NewAuthenticationTokenFromFile
NewAuthenticationTokenFromFile return a interface of a Provider with a string token file path.
pulsar/auth/token.go:68
↓ 2 callers
Function
NewAuthenticationTokenFromSupplier
NewAuthenticationTokenFromSupplier returns a token auth provider that gets the token data from a user supplied function. The function is invoked each
pulsar/client.go:51
↓ 2 callers
Function
NewBufferPool
()
pulsar/internal/buffer.go:97
↓ 2 callers
Function
NewBundlesDataWithNumBundles
(numBundles int)
pulsaradmin/pkg/utils/bundles_data.go:32
↓ 2 callers
Function
NewClient
()
perf/pulsar-perf-go.go:53
↓ 2 callers
Function
NewDefaultBackoffWithInitialBackOff
(backoff time.Duration)
pulsar/backoff/backoff.go:50
↓ 2 callers
Function
NewDefaultTokenCache
(audience string, flow *oauth2.ClientCredentialsFlow)
oauth2/cache/cache.go:53
↓ 2 callers
Function
NewMessageReader
(headersAndPayload Buffer)
pulsar/internal/commands.go:54
← previous
next →
601–700 of 5,094, ranked by callers