Metadata-Version: 2.4
Name: zemu-kafka-sdk
Version: 1.0.0
Summary: zemu standard Kafka message SDK
Project-URL: Homepage, https://github.com/TFDatas
Project-URL: Documentation, https://github.com/TFDatas/kafka/blob/master/README.md
Project-URL: Source, https://github.com/TFDatas/kafka
Author-email: zemu Team <tfkxsx@gmail.com>
Maintainer-email: zemu Team <tfkxsx@gmail.com>
License-Expression: MIT
Keywords: data-pipeline,distributed-tracing,kafka,message-queue,sdk,zemu
Classifier: Development Status :: 4 - Beta
Classifier: Intended Audience :: Developers
Classifier: License :: OSI Approved :: MIT License
Classifier: Operating System :: OS Independent
Classifier: Programming Language :: Python :: 3
Classifier: Programming Language :: Python :: 3.9
Classifier: Programming Language :: Python :: 3.10
Classifier: Programming Language :: Python :: 3.11
Classifier: Programming Language :: Python :: 3.12
Classifier: Programming Language :: Python :: 3.13
Classifier: Topic :: Communications
Classifier: Topic :: Software Development :: Libraries :: Python Modules
Classifier: Topic :: System :: Distributed Computing
Requires-Python: >=3.9
Requires-Dist: confluent-kafka>=2.0.2
Provides-Extra: all
Requires-Dist: avro>=1.11.0; extra == 'all'
Requires-Dist: requests>=2.28.0; extra == 'all'
Provides-Extra: avro
Requires-Dist: avro>=1.11.0; extra == 'avro'
Provides-Extra: dev
Requires-Dist: black>=23.0.0; extra == 'dev'
Requires-Dist: build>=1.2.0; extra == 'dev'
Requires-Dist: flake8>=6.0.0; extra == 'dev'
Requires-Dist: isort>=5.12.0; extra == 'dev'
Requires-Dist: mypy>=1.0.0; extra == 'dev'
Requires-Dist: pytest-cov>=4.0.0; extra == 'dev'
Requires-Dist: pytest>=7.0.0; extra == 'dev'
Requires-Dist: twine>=5.0.0; extra == 'dev'
Provides-Extra: schema-registry
Requires-Dist: requests>=2.28.0; extra == 'schema-registry'
Description-Content-Type: text/markdown

# kafka-sdk

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

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

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

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

## 1. 安装

本地源码开发：

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

业务项目源码引入：

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

构建发布包：

```bash
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
```

生成产物：

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

## 2. 依赖

基础依赖：

```text
confluent-kafka>=2.0.2
```

默认推荐使用 JSON 序列化：

```bash
export KAFKA_SERIALIZATION_FORMAT=json
```

如果要使用 Avro：

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

如果要调用 Schema Registry HTTP API：

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

开发依赖：

```bash
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：

```json
{
  "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` |

完整映射见：

```text
kafkaSdk/topic_registry.py
```

## 7. Producer 示例

```python
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：

```python
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：

```python
from kafkaSdk import MessageHandler, ProcessingResult


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

## 9. 任务上下文校验

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

```python
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：

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

在 Python 中读取 Topic 定义：

```python
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：

```bash
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 文件：

```text
kafkaSdk/schema/standard_message.json
```

## 12. 发布前检查

```bash
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 内容：

```bash
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`
