brk-code

kafka-python の API リファレンス

kafka-python (dpkp/kafka-python) の公開 API 89 件 —— クラス 64、関数 6、メソッド 19。実際のソースを静的解析して抽出した正確なシグネチャを掲載しています。

リポジトリ: dpkp/kafka-python

種別件数
クラス64
関数6
メソッド19

API 一覧

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.

この情報について

掲載しているシグネチャは dpkp/kafka-python の公開ソースコードを Python の ast モジュールで静的解析し、引数名・デフォルト値・ 型注釈・戻り値型をそのまま抽出したものです。実装コードは保存していません。 詳しくは仕組みの解説をご覧ください。

収録ライブラリ一覧(全 805 件)へ戻る