Skip to main content

kafka-sdk

kafka-sdk 是 tf-data 系统的标准 Kafka 消息 SDK,Python 导入包名为 kafkaSdk

SDK 只支持当前冻结的最终消息协议:

  • 标准 envelope
  • 标准 message_type
  • 标准 Topic 路由
  • 标准任务上下文校验
  • 标准 producer / consumer 入口

它不再兼容历史 Topic alias、旧 message_type、旧 handler 命名或过渡协议。

1. 安装

本地源码开发:

cd /Users/tf/Documents/worke/tfd/kafka
/Users/tf/Documents/worke/tfd/api/.venv/bin/python -m pip install -e .

业务项目源码引入:

export TFD_ROOT=/Users/tf/Documents/worke/tfd
export PYTHONPATH="$TFD_ROOT/kafka:$TFD_ROOT/logging-sdk/src:$PWD"

构建发布包:

cd /Users/tf/Documents/worke/tfd/kafka
/Users/tf/Documents/worke/tfd/api/.venv/bin/python -m pip install build twine
/Users/tf/Documents/worke/tfd/api/.venv/bin/python -m build

生成产物:

dist/kafka_sdk-<version>-py3-none-any.whl
dist/kafka_sdk-<version>.tar.gz

2. 依赖

基础依赖:

confluent-kafka>=2.0.2

默认推荐使用 JSON 序列化:

export KAFKA_SERIALIZATION_FORMAT=json

如果要使用 Avro:

pip install "kafka-sdk[avro]"
export KAFKA_SERIALIZATION_FORMAT=avro
export SCHEMA_REGISTRY_URL=http://192.168.2.10:8081

如果要调用 Schema Registry HTTP API:

pip install "kafka-sdk[schema-registry]"

开发依赖:

pip install "kafka-sdk[dev]"

3. 环境变量

变量 默认值 说明
BOOTSTRAP_SERVERS localhost:9092 Kafka bootstrap servers
SCHEMA_REGISTRY_URL http://localhost:8081 Avro Schema Registry
KAFKA_SERIALIZATION_FORMAT avro jsonavro,业务当前推荐 json
KAFKA_SECURITY_PROTOCOL PLAINTEXT Kafka 安全协议
KAFKA_SASL_MECHANISM PLAIN SASL 机制
KAFKA_SASL_USERNAME SASL 用户名
KAFKA_SASL_PASSWORD SASL 密码
KAFKA_PRODUCER_CLIENT_ID tf-data-producer Producer client id
KAFKA_CONSUMER_GROUP_ID tf-data-consumer-group Consumer group id
KAFKA_CONSUMER_CLIENT_ID tf-data-consumer Consumer client id
KAFKA_CONSUMER_AUTO_OFFSET_RESET earliest offset reset 策略
KAFKA_CONSUMER_ENABLE_AUTO_COMMIT false 是否自动提交 offset
DEPLOYMENT_REGION both internaloverseasbothglobal
KAFKA_TOPIC_DLQ_TASK dlq.task 死信 Topic

4. 标准 Envelope

所有消息必须使用标准 envelope:

{
  "trace_id": "string",
  "task_id": "string",
  "timestamp": "2026-05-06T10:00:00.000000+00:00",
  "payload": {},
  "metadata": {
    "message_type": "cmd.clean",
    "schema_version": "v1",
    "producer": "tf-data",
    "region": "internal",
    "retry_count": 0,
    "scope": {
      "tenant_id": "tenant-001",
      "project_id": "project-001"
    }
  }
}

5. 标准 Topic

命令 Topic:

Topic 用途
cmd.download.internal 国内下载命令
cmd.download.overseas 海外下载命令
cmd.extract 正式抽取命令
cmd.extract.online_test 在线测试抽取命令
cmd.clean 清洗命令
cmd.generate_file 文件生成/导出命令
cmd.autoparse AutoParse 命令
cmd.email.send 加密的最终 Email 发送命令,固定 1 个分区

状态、结果和系统 Topic:

Topic 用途
status.task 统一任务状态
status.download 下载阶段状态
status.extract 抽取阶段状态
result.extract 结果摘要
event.system 系统事件
dlq.task 死信任务
status.email.delivery SMTP 单次发送结果
dlq.email.delivery 非法或最终失败的 Email 投递死信

6. message_type 路由

SDK 使用 message_type 路由到物理 Topic。

示例:

message_type Topic
cmd.download.internal cmd.download.internal
cmd.download.overseas cmd.download.overseas
cmd.online_test.download.internal cmd.download.internal
cmd.online_test.download.overseas cmd.download.overseas
cmd.online_test.extract cmd.extract.online_test
result.clean result.extract
result.generate_file result.extract
result.autoparse result.extract
result.online_test result.extract
cmd.email.send cmd.email.send
status.email.delivery status.email.delivery
dlq.email.delivery dlq.email.delivery

完整映射见:

kafkaSdk/topic_registry.py

7. Producer 示例

from kafkaSdk import KafkaConfig, KafkaProducer

config = KafkaConfig(
    bootstrap_servers="192.168.2.10:19092",
    serialization_format="json",
)
producer = KafkaProducer(config=config, producer_id="example-producer")

producer.send_by_message_type(
    message_type="cmd.clean",
    payload={
        "task_context": {
            "task_id": "task-001",
            "tenant_id": "tenant-001",
            "project_id": "project-001",
            "task_type": "CLEAN",
            "schedule_type": "MANUAL_TASK",
            "request_source": "readme-example",
            "created_by": "developer",
            "operator": "developer",
            "business_payload": {"table": "result_demo"},
        }
    },
    task_id="task-001",
    scope={"tenant_id": "tenant-001", "project_id": "project-001"},
    region="internal",
    synchronous=True,
)

producer.close()

8. Consumer 示例

函数 handler:

from kafkaSdk import KafkaConfig, KafkaConsumer, ProcessingResult


def handle_clean(payload, metadata, context):
    task_id = payload.get("task_context", {}).get("task_id")
    print(f"clean task: {task_id}")
    return ProcessingResult.SUCCESS


config = KafkaConfig(
    bootstrap_servers="192.168.2.10:19092",
    serialization_format="json",
    consumer_group_id="clean-worker",
)

consumer = KafkaConsumer(config=config)
consumer.register_message_handler("cmd.clean", handle_clean)
consumer.subscribe(["cmd.clean"])
consumer.start()

类 handler:

from kafkaSdk import MessageHandler, ProcessingResult


class CleanHandler(MessageHandler):
    def handle(self, payload, metadata, context):
        return ProcessingResult.SUCCESS

9. 任务上下文校验

SDK 提供任务上下文契约,用于在 producer 或 handler 入口校验关键字段。

from kafkaSdk import TaskCategory, validate_task_command_context

context = {
    "task_id": "task-001",
    "parent_task_id": "task-001",
    "tenant_id": "tenant-001",
    "project_id": "project-001",
    "task_type": "CLEAN",
    "schedule_type": "MANUAL_TASK",
    "request_source": "readme-example",
    "created_by": "developer",
    "operator": "developer",
    "business_payload": {"table": "result_demo"},
}

result = validate_task_command_context(TaskCategory.CLEAN, context)
if not result.valid:
    raise ValueError(
        f"invalid context, missing={result.missing_fields}, empty={result.empty_fields}"
    )

10. Topic 初始化

当前统一部署入口是 deploy/,SDK README 不再提供历史部署包命令。

在服务器上创建最终 Topic:

cd /home/tf/tfd/deploy/scripts
./create_final_topics_with_docker.sh kafka:9092

在 Python 中读取 Topic 定义:

from kafkaSdk import list_topic_definitions

for definition in list_topic_definitions():
    print(definition.name, definition.partitions, definition.description)

11. Schema 注册

当前本地默认 JSON 序列化,可不注册 Avro Schema。

如果切换 Avro:

cd /home/tf/tfd/deploy/scripts
./register_standard_schema.sh \
  http://192.168.2.10:8081 \
  /home/tf/tfd/kafka/kafkaSdk/schema/standard_message.json

Schema 文件:

kafkaSdk/schema/standard_message.json

12. 发布前检查

cd /Users/tf/Documents/worke/tfd/kafka
/Users/tf/Documents/worke/tfd/api/.venv/bin/python -m py_compile kafkaSdk/*.py
/Users/tf/Documents/worke/tfd/api/.venv/bin/python -m build

检查 wheel 内容:

python - <<'PY'
from pathlib import Path
from zipfile import ZipFile

wheel = sorted(Path("dist").glob("*.whl"))[-1]
with ZipFile(wheel) as zf:
    for name in zf.namelist():
        print(name)
PY

13. 净化门禁

发布前必须满足:

  1. 不解析历史 Topic alias。
  2. 不路由旧 message_type
  3. 不注册旧 handler。
  4. Producer 发送未知 message_type 时直接报错。
  5. Consumer 收到未知 message_type 时进入失败或死信路径。
  6. Topic 初始化统一走 deploy/scripts/create_final_topics_with_docker.sh
  7. README 不再提供旧部署包命令。

14. 版本

当前版本:1.0.0

版本号需要同时保持一致:

  • pyproject.toml
  • kafkaSdk/__init__.py

Release files for zemu-kafka-sdk 1.0.1

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for zemu-kafka-sdk 1.0.1
File Size Uploaded
zemu_kafka_sdk-1.0.1.tar.gz 41.4 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for zemu-kafka-sdk 1.0.1
File Interpreter ABI Platform
zemu_kafka_sdk-1.0.1-py3-none-any.whl Python 3 none any Details

Total release size: 88.2 kB

Release files / zemu_kafka_sdk-1.0.1.tar.gz

Download URL zemu_kafka_sdk-1.0.1.tar.gz
Size 41.4 kB
Tags Source
SHA-256 checksum
How to use checksums
191b2e1d6bedad0fce4c95c9db34aaf2e85056ae767da6b62e7736a7480d8456
BLAKE2b-256 checksum
How to use checksums
7e3ae795167ab093543f002df55306266b1039fd514434b62cb61fbf64a65177
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/7.0.0 CPython/3.12.0

Release files / zemu_kafka_sdk-1.0.1-py3-none-any.whl

Download URL zemu_kafka_sdk-1.0.1-py3-none-any.whl
Size 46.7 kB
Tags Python 3
SHA-256 checksum
How to use checksums
fc856ba3e0385f363af9c0bd01feb35732ccf5b8fd8ca822c73779f4ef59cb64
BLAKE2b-256 checksum
How to use checksums
357423c3bf6e72fda961dcb0b9a72fa7ef824ea39c4d0377826c144d0858b08c
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/7.0.0 CPython/3.12.0

Release history Release notifications | RSS feed

This release

1.0.1 This release

2 release files

1.0.0

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page