kafka-python API reference
89 public APIs from kafka-python (dpkp/kafka-python) — 64 classes, 6 functions, 19 methods. Signatures extracted by static analysis of the actual source.
Repository: dpkp/kafka-python
| Kind | Count |
|---|---|
| Classes | 64 |
| Functions | 6 |
| Methods | 19 |
API list
class
kafka.admin._acls.ACLRepresents a concrete ACL for a specific ResourcePattern.
class
kafka.admin._acls.ACLAdminMixinMixin providing ACL management methods for KafkaAdminClient.
class
kafka.admin._acls.ACLFilterRepresents a filter to use with describing and deleting ACLs.
class
kafka.admin._acls.ACLOperationType of operation.
class
kafka.admin._acls.ACLPermissionTypeAn enumerated type of permissions.
class
kafka.admin._acls.ACLResourcePatternTypeAn enumerated type of resource patterns.
class
kafka.admin._acls.ResourcePatternA resource pattern to apply the ACL to.
class
kafka.admin._acls.ResourceTypeType of kafka resource to set ACL for.
class
kafka.admin._configs.ConfigResourceA class for specifying config resources.
class
kafka.admin._groups.GroupTypeConsumer group protocol types (KIP-848).
class
kafka.admin._partitions.PartitionAdminMixinMixin providing partition and record management methods.
class
kafka.admin._topics.NewTopicDEPRECATED: A class for new topic creation.
class
kafka.admin._transactions.AbortTransactionSpecInputs for ``abort_transaction``.
class
kafka.admin._transactions.ProducerStateOne ActiveProducer row from DescribeProducers.
class
kafka.admin._transactions.TransactionListingOne row from a ListTransactions response.
class
kafka.admin._transactions.TransactionsAdminMixinMixin providing KIP-664 hanging-transaction tooling.
class
kafka.admin._users.UserAdminMixinMixin providing user management methods for KafkaAdminClient.
class
kafka.admin._users.UserScramCredentialDeletionSpecifies that a SCRAM credential should be deleted.
class
kafka.admin.client.KafkaAdminClientA class for administering the Kafka cluster.
class
kafka.cluster.ClusterMetadataA class to manage kafka cluster metadata.
method
kafka.cluster.ClusterMetadata.topics(exclude_internal_topics=True)Get set of known topics.
class
kafka.consumer.group.KafkaConsumerConsume records from a Kafka cluster.
method
kafka.consumer.group.KafkaConsumer.subscription()Get the current topic subscription.
class
kafka.coordinator.assignors.range.RangePartitionAssignorThe range assignor works on a per-topic basis.
class
kafka.errors.IncompatibleBrokerVersionSynthetic error raised by client
class
kafka.metrics.dict_reporter.DictReporterA basic dictionary based metrics reporter.
class
kafka.metrics.measurable.AbstractMeasurableA measurable quantity that can be registered as a metric
class
kafka.metrics.metric_config.MetricConfigConfiguration values for metrics
class
kafka.metrics.metrics.MetricsA registry of sensors and metrics.
method
kafka.metrics.metrics.Metrics.add_reporter(reporter)Add a MetricReporter
method
kafka.metrics.metrics.Metrics.close()Close this metrics repository.
class
kafka.metrics.quota.QuotaAn upper or lower bound for metrics
class
kafka.metrics.stats.max_stat.MaxAn AbstractSampledStat that gives the max over its samples.
class
kafka.metrics.stats.min_stat.MinAn AbstractSampledStat that gives the min over its samples.
class
kafka.metrics.stats.percentiles.PercentilesA compound stat that reports one or more percentiles
class
kafka.metrics.stats.rate.RateThe rate of the given quantity.
class
kafka.metrics.stats.total.TotalAn un-windowed cumulative total maintained over all time.
class
kafka.net.backend.abstract.NetBackendContract for a pluggable async event-loop backend.
method
kafka.net.backend.abstract.NetBackend.await_for(future:Any, timeout_ms:Optional[float], raise_error:bool=True) -> AnyAwait ``future`` with a timeout in ms.
method
kafka.net.backend.abstract.NetBackend.call_at(when:float, task:Any) -> AnySchedule ``task`` to run at absolute monotonic time ``when``.
method
kafka.net.backend.abstract.NetBackend.call_later(delay:float, task:Any) -> AnySchedule ``task`` to run after ``delay`` seconds.
method
kafka.net.backend.abstract.NetBackend.call_soon(task:Any) -> AnyEnqueue a coroutine/callable to run on the next loop iteration.
method
kafka.net.backend.abstract.NetBackend.cancel(task:Any) -> NoneCancel a scheduled task/timer previously returned by call_*.
method
kafka.net.backend.abstract.NetBackend.close() -> NoneStop (if running) and release loop resources.
method
kafka.net.backend.abstract.NetBackend.getaddrinfo(host:str, port:int) -> AddrInfoResultResolve host/port via DNS
method
kafka.net.backend.abstract.NetBackend.sleep(delay:float) -> AnyAwaitable that resolves after ``delay`` seconds.
method
kafka.net.backend.abstract.NetBackend.start() -> NoneSpawn/attach the IO thread that runs the loop.
method
kafka.net.backend.abstract.NetBackend.stop(timeout_ms:Optional[float]=None) -> NoneStop the loop and join the IO thread.
method
kafka.net.backend.abstract.NetBackend.wakeup() -> NoneInterrupt the loop's select() from another thread.
class
kafka.net.backend.abstract.NetTransportThe transport surface used by NetProtocol / NetBackend.
func
kafka.net.backend.abstract.list_backends()List all registered backends.
class
kafka.net.backend.asyncio_backend.AsyncioFuture``create_future()`` result for the asyncio backend.
class
kafka.net.socks5.Socks5ProxyProtocolSocks5 proxy sans-IO protocol handler.
class
kafka.partitioner.abc.PartitionerBase class for pluggable partition selection strategies.
class
kafka.partitioner.default.DefaultPartitionerDefault partitioner.
func
kafka.partitioner.default.murmur2(data)Pure-python Murmur2 implementation.
class
kafka.producer.kafka.KafkaProducerA Kafka client that publishes records to the Kafka cluster.
method
kafka.producer.kafka.KafkaProducer.abort_transaction()Aborts the ongoing transaction.
method
kafka.producer.kafka.KafkaProducer.close(timeout=None, null_logger=False)Close this producer.
method
kafka.producer.kafka.KafkaProducer.commit_transaction()Commits the ongoing transaction.
class
kafka.producer.sender.SenderDrives the sending of produce requests to the Kafka cluster.
method
kafka.producer.sender.Sender.wakeup()Wake the sender loop early (e.g.
class
kafka.producer.transaction_manager.TransactionManagerA class which maintains state for transactions.
class
kafka.protocol.admin.topics.ElectionTypeLeader election type
class
kafka.protocol.consumer.offsets.OffsetTimestampMillisecond-timestamp spec for partition offset lookup.
func
kafka.protocol.generate_stubs.generate_all(dry_run=False, check=False)Generate all stub files.
class
kafka.protocol.old.admin.DescribeAclsRequest_v2Enable flexible version
class
kafka.protocol.old.admin.ElectionTypeLeader election type
class
kafka.protocol.old.fetch.FetchRequest_v10bumped up to indicate ZStandard capability.
class
kafka.protocol.old.fetch.FetchRequest_v11added rack ID to support read from followers (KIP-392)
class
kafka.protocol.old.fetch.FetchRequest_v7Add incremental fetch requests (see KIP-227)
class
kafka.protocol.old.fetch.FetchRequest_v9adds the current leader epoch (see KIP-320)
class
kafka.protocol.old.fetch.FetchResponse_v6Same as FetchResponse_v5.
class
kafka.protocol.old.fetch.FetchResponse_v7Add error_code and session_id to response
class
kafka.protocol.old.list_offsets.ListOffsetsRequest_v4Add current_leader_epoch to request
class
kafka.protocol.old.list_offsets.ListOffsetsResponse_v4Add leader_epoch to response
class
kafka.protocol.old.list_offsets.ListOffsetsResponse_v5adds a new error code, OFFSET_NOT_AVAILABLE
class
kafka.protocol.old.metadata.MetadataRequest_v5The v5 metadata request is the same as v4.
class
kafka.protocol.old.metadata.MetadataResponse_v7v7 adds per-partition leader_epoch field
class
kafka.protocol.old.metadata.MetadataResponse_v8v8 adds authorized_operations fields
class
kafka.protocol.old.produce.ProduceRequest_v5Same as v4.
class
kafka.protocol.old.produce.ProduceRequest_v7V7 bumped up to indicate ZStandard capability.
class
kafka.protocol.old.produce.ProduceResponse_v7V7 bumped up to indicate ZStandard capability.
class
kafka.protocol.sasl.SaslBytesResponseResponse for raw SASL v0 exchange -- returns bytes as-is.
class
kafka.protocol.schemas.fields.codecs.types.FixedCodecBase class for fixed-size codecs.
class
kafka.protocol.schemas.fields.codegen.CodegenContextShared state for code generation.
func
kafka.record._crc32c.crc(data)Compute CRC-32C checksum of the data.
func
kafka.record._crc32c.crc_finalize(crc)Finalize CRC-32C checksum.
func
kafka.record._crc32c.crc_update(crc, data, _TABLE=CRC_TABLE, _M=_MASK)Update CRC-32C checksum with data.
About this data
These signatures were extracted from the public source of dpkp/kafka-python
using Python's ast module. Argument names, default values,
type annotations and return types are taken verbatim from the code.
Implementation bodies are never stored. See
how it works for details.