kafka-python の API リファレンス
kafka-python (dpkp/kafka-python) の公開 API 89 件 —— クラス 64、関数 6、メソッド 19。実際のソースを静的解析して抽出した正確なシグネチャを掲載しています。
リポジトリ: dpkp/kafka-python
| 種別 | 件数 |
|---|---|
| クラス | 64 |
| 関数 | 6 |
| メソッド | 19 |
API 一覧
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.
この情報について
掲載しているシグネチャは dpkp/kafka-python の公開ソースコードを
Python の ast モジュールで静的解析し、引数名・デフォルト値・
型注釈・戻り値型をそのまま抽出したものです。実装コードは保存していません。
詳しくは仕組みの解説をご覧ください。