sdkagent

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

KindCount
Classes64
Functions6
Methods19

API list

classkafka.admin._acls.ACL
Represents a concrete ACL for a specific ResourcePattern.
classkafka.admin._acls.ACLAdminMixin
Mixin providing ACL management methods for KafkaAdminClient.
classkafka.admin._acls.ACLFilter
Represents a filter to use with describing and deleting ACLs.
classkafka.admin._acls.ACLOperation
Type of operation.
classkafka.admin._acls.ACLPermissionType
An enumerated type of permissions.
classkafka.admin._acls.ACLResourcePatternType
An enumerated type of resource patterns.
classkafka.admin._acls.ResourcePattern
A resource pattern to apply the ACL to.
classkafka.admin._acls.ResourceType
Type of kafka resource to set ACL for.
classkafka.admin._configs.ConfigResource
A class for specifying config resources.
classkafka.admin._groups.GroupType
Consumer group protocol types (KIP-848).
classkafka.admin._partitions.PartitionAdminMixin
Mixin providing partition and record management methods.
classkafka.admin._topics.NewTopic
DEPRECATED: A class for new topic creation.
classkafka.admin._transactions.AbortTransactionSpec
Inputs for ``abort_transaction``.
classkafka.admin._transactions.ProducerState
One ActiveProducer row from DescribeProducers.
classkafka.admin._transactions.TransactionListing
One row from a ListTransactions response.
classkafka.admin._transactions.TransactionsAdminMixin
Mixin providing KIP-664 hanging-transaction tooling.
classkafka.admin._users.UserAdminMixin
Mixin providing user management methods for KafkaAdminClient.
classkafka.admin._users.UserScramCredentialDeletion
Specifies that a SCRAM credential should be deleted.
classkafka.admin.client.KafkaAdminClient
A class for administering the Kafka cluster.
classkafka.cluster.ClusterMetadata
A class to manage kafka cluster metadata.
methodkafka.cluster.ClusterMetadata.topics(exclude_internal_topics=True)
Get set of known topics.
classkafka.consumer.group.KafkaConsumer
Consume records from a Kafka cluster.
methodkafka.consumer.group.KafkaConsumer.subscription()
Get the current topic subscription.
classkafka.coordinator.assignors.range.RangePartitionAssignor
The range assignor works on a per-topic basis.
classkafka.errors.IncompatibleBrokerVersion
Synthetic error raised by client
classkafka.metrics.dict_reporter.DictReporter
A basic dictionary based metrics reporter.
classkafka.metrics.measurable.AbstractMeasurable
A measurable quantity that can be registered as a metric
classkafka.metrics.metric_config.MetricConfig
Configuration values for metrics
classkafka.metrics.metrics.Metrics
A registry of sensors and metrics.
methodkafka.metrics.metrics.Metrics.add_reporter(reporter)
Add a MetricReporter
methodkafka.metrics.metrics.Metrics.close()
Close this metrics repository.
classkafka.metrics.quota.Quota
An upper or lower bound for metrics
classkafka.metrics.stats.max_stat.Max
An AbstractSampledStat that gives the max over its samples.
classkafka.metrics.stats.min_stat.Min
An AbstractSampledStat that gives the min over its samples.
classkafka.metrics.stats.percentiles.Percentiles
A compound stat that reports one or more percentiles
classkafka.metrics.stats.rate.Rate
The rate of the given quantity.
classkafka.metrics.stats.total.Total
An un-windowed cumulative total maintained over all time.
classkafka.net.backend.abstract.NetBackend
Contract for a pluggable async event-loop backend.
methodkafka.net.backend.abstract.NetBackend.await_for(future:Any, timeout_ms:Optional[float], raise_error:bool=True) -> Any
Await ``future`` with a timeout in ms.
methodkafka.net.backend.abstract.NetBackend.call_at(when:float, task:Any) -> Any
Schedule ``task`` to run at absolute monotonic time ``when``.
methodkafka.net.backend.abstract.NetBackend.call_later(delay:float, task:Any) -> Any
Schedule ``task`` to run after ``delay`` seconds.
methodkafka.net.backend.abstract.NetBackend.call_soon(task:Any) -> Any
Enqueue a coroutine/callable to run on the next loop iteration.
methodkafka.net.backend.abstract.NetBackend.cancel(task:Any) -> None
Cancel a scheduled task/timer previously returned by call_*.
methodkafka.net.backend.abstract.NetBackend.close() -> None
Stop (if running) and release loop resources.
methodkafka.net.backend.abstract.NetBackend.getaddrinfo(host:str, port:int) -> AddrInfoResult
Resolve host/port via DNS
methodkafka.net.backend.abstract.NetBackend.sleep(delay:float) -> Any
Awaitable that resolves after ``delay`` seconds.
methodkafka.net.backend.abstract.NetBackend.start() -> None
Spawn/attach the IO thread that runs the loop.
methodkafka.net.backend.abstract.NetBackend.stop(timeout_ms:Optional[float]=None) -> None
Stop the loop and join the IO thread.
methodkafka.net.backend.abstract.NetBackend.wakeup() -> None
Interrupt the loop's select() from another thread.
classkafka.net.backend.abstract.NetTransport
The transport surface used by NetProtocol / NetBackend.
funckafka.net.backend.abstract.list_backends()
List all registered backends.
classkafka.net.backend.asyncio_backend.AsyncioFuture
``create_future()`` result for the asyncio backend.
classkafka.net.socks5.Socks5ProxyProtocol
Socks5 proxy sans-IO protocol handler.
classkafka.partitioner.abc.Partitioner
Base class for pluggable partition selection strategies.
classkafka.partitioner.default.DefaultPartitioner
Default partitioner.
funckafka.partitioner.default.murmur2(data)
Pure-python Murmur2 implementation.
classkafka.producer.kafka.KafkaProducer
A Kafka client that publishes records to the Kafka cluster.
methodkafka.producer.kafka.KafkaProducer.abort_transaction()
Aborts the ongoing transaction.
methodkafka.producer.kafka.KafkaProducer.close(timeout=None, null_logger=False)
Close this producer.
methodkafka.producer.kafka.KafkaProducer.commit_transaction()
Commits the ongoing transaction.
classkafka.producer.sender.Sender
Drives the sending of produce requests to the Kafka cluster.
methodkafka.producer.sender.Sender.wakeup()
Wake the sender loop early (e.g.
classkafka.producer.transaction_manager.TransactionManager
A class which maintains state for transactions.
classkafka.protocol.admin.topics.ElectionType
Leader election type
classkafka.protocol.consumer.offsets.OffsetTimestamp
Millisecond-timestamp spec for partition offset lookup.
funckafka.protocol.generate_stubs.generate_all(dry_run=False, check=False)
Generate all stub files.
classkafka.protocol.old.admin.DescribeAclsRequest_v2
Enable flexible version
classkafka.protocol.old.admin.ElectionType
Leader election type
classkafka.protocol.old.fetch.FetchRequest_v10
bumped up to indicate ZStandard capability.
classkafka.protocol.old.fetch.FetchRequest_v11
added rack ID to support read from followers (KIP-392)
classkafka.protocol.old.fetch.FetchRequest_v7
Add incremental fetch requests (see KIP-227)
classkafka.protocol.old.fetch.FetchRequest_v9
adds the current leader epoch (see KIP-320)
classkafka.protocol.old.fetch.FetchResponse_v6
Same as FetchResponse_v5.
classkafka.protocol.old.fetch.FetchResponse_v7
Add error_code and session_id to response
classkafka.protocol.old.list_offsets.ListOffsetsRequest_v4
Add current_leader_epoch to request
classkafka.protocol.old.list_offsets.ListOffsetsResponse_v4
Add leader_epoch to response
classkafka.protocol.old.list_offsets.ListOffsetsResponse_v5
adds a new error code, OFFSET_NOT_AVAILABLE
classkafka.protocol.old.metadata.MetadataRequest_v5
The v5 metadata request is the same as v4.
classkafka.protocol.old.metadata.MetadataResponse_v7
v7 adds per-partition leader_epoch field
classkafka.protocol.old.metadata.MetadataResponse_v8
v8 adds authorized_operations fields
classkafka.protocol.old.produce.ProduceRequest_v5
Same as v4.
classkafka.protocol.old.produce.ProduceRequest_v7
V7 bumped up to indicate ZStandard capability.
classkafka.protocol.old.produce.ProduceResponse_v7
V7 bumped up to indicate ZStandard capability.
classkafka.protocol.sasl.SaslBytesResponse
Response for raw SASL v0 exchange -- returns bytes as-is.
classkafka.protocol.schemas.fields.codecs.types.FixedCodec
Base class for fixed-size codecs.
classkafka.protocol.schemas.fields.codegen.CodegenContext
Shared state for code generation.
funckafka.record._crc32c.crc(data)
Compute CRC-32C checksum of the data.
funckafka.record._crc32c.crc_finalize(crc)
Finalize CRC-32C checksum.
funckafka.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.

Back to all 805 libraries