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 |
json 或 avro,业务当前推荐 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 |
internal、overseas、both、global |
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. 净化门禁
发布前必须满足:
- 不解析历史 Topic alias。
- 不路由旧
message_type。 - 不注册旧 handler。
- Producer 发送未知
message_type时直接报错。 - Consumer 收到未知
message_type时进入失败或死信路径。 - Topic 初始化统一走
deploy/scripts/create_final_topics_with_docker.sh。 - README 不再提供旧部署包命令。
14. 版本
当前版本:1.0.0
版本号需要同时保持一致:
pyproject.tomlkafkaSdk/__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)
| File | Size | Uploaded | |
|---|---|---|---|
| zemu_kafka_sdk-1.0.1.tar.gz | 41.4 kB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| 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
|