架构核对发现 2 处差距,本次补齐: - Kafka 上行支持 SSL/SASL_SSL 双向 mTLS(8.2 服务间 mTLS 边缘网关↔总线) - 网关日志支持结构化 JSON 输出(NFR 可维护:统一日志规范) - 配置示例与 README 验收口径同步更新
160 lines
6.4 KiB
Python
160 lines
6.4 KiB
Python
# -*- coding: utf-8 -*-
|
|
"""Kafka 流式上行 —— 模板化封装(topic 命名随模板变化)+ 断点续传确认。
|
|
|
|
设计:
|
|
- topic 模板:`{topic_prefix}.{device_id}.points`(如 `ti-cl4.CLF-01.points`);
|
|
- 每条样本先写 spool(engine 已写),发送成功后 ack 删除;
|
|
- Kafka 不可用时:生产者切换为降级模式,样本保留在 spool,
|
|
下次 publish 时从 pending 重发(断点续传,丢失率受保障);
|
|
- 依赖 `confluent-kafka`;未安装时自动降级为 NullProducer(本地联调用)。
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
import threading
|
|
import time
|
|
from typing import List, Optional
|
|
|
|
from collector.spool import SpoolStore
|
|
|
|
logger = logging.getLogger("edge_gateway.kafka_sink")
|
|
|
|
|
|
class KafkaSink:
|
|
"""Kafka 上行通道。"""
|
|
|
|
def __init__(
|
|
self,
|
|
bootstrap_servers: str,
|
|
topic_prefix: str,
|
|
spool: SpoolStore,
|
|
batch_size: int = 500,
|
|
security: Optional[dict] = None,
|
|
):
|
|
"""
|
|
Args:
|
|
bootstrap_servers: Kafka 地址列表(逗号分隔);
|
|
topic_prefix: 上行 topic 前缀(模板化 `{prefix}.{device}.points`);
|
|
spool: 本地缓存(断点续传);
|
|
batch_size: 单批上限;
|
|
security: 传输安全配置(NFR 9 章:服务间 mTLS)。
|
|
约定字段:protocol(PLAINTEXT|SSL|SASL_SSL)、ca_location、
|
|
cert_location、key_location、sasl_username、sasl_password。
|
|
其中 SSL/SASL_SSL + ca/cert/key = 边缘网关↔总线双向 mTLS。
|
|
"""
|
|
self.bootstrap_servers = bootstrap_servers
|
|
self.topic_prefix = topic_prefix
|
|
self.spool = spool
|
|
self.batch_size = max(1, batch_size)
|
|
self.security = security or {}
|
|
self._producer = None
|
|
self._degraded = False # True = Kafka 不可用,仅保留 spool
|
|
self._lock = threading.Lock()
|
|
self._sent = 0
|
|
self._failed = 0
|
|
self._connect()
|
|
|
|
# ------------------------------------------------------------------
|
|
def _producer_config(self) -> dict:
|
|
"""组装 confluent-kafka producer 配置(含 mTLS/SASL 透传)。"""
|
|
conf = {"bootstrap.servers": self.bootstrap_servers}
|
|
protocol = self.security.get("protocol", "PLAINTEXT").upper()
|
|
conf["security.protocol"] = protocol
|
|
if protocol in ("SSL", "SASL_SSL"):
|
|
for key, kafka_key in (
|
|
("ca_location", "ssl.ca.location"),
|
|
("cert_location", "ssl.certificate.location"),
|
|
("key_location", "ssl.key.location"),
|
|
("key_password", "ssl.key.password"),
|
|
):
|
|
if self.security.get(key):
|
|
conf[kafka_key] = self.security[key]
|
|
if protocol == "SASL_SSL":
|
|
conf["sasl.mechanism"] = self.security.get("sasl_mechanism", "PLAIN")
|
|
if self.security.get("sasl_username"):
|
|
conf["sasl.username"] = self.security["sasl_username"]
|
|
conf["sasl.password"] = self.security.get("sasl_password", "")
|
|
return conf
|
|
|
|
def _connect(self) -> None:
|
|
"""尝试连接 Kafka;失败则降级(不阻断采集)。"""
|
|
try:
|
|
from confluent_kafka import Producer # type: ignore
|
|
except ImportError:
|
|
logger.warning("confluent-kafka 未安装,Kafka 上行降级为 spool-only 模式")
|
|
self._degraded = True
|
|
return
|
|
try:
|
|
self._producer = Producer(self._producer_config())
|
|
except Exception as exc:
|
|
logger.warning("Kafka 初始化失败(%s),降级为 spool-only 模式", exc)
|
|
self._degraded = True
|
|
|
|
def _topic(self, device_id: str) -> str:
|
|
return f"{self.topic_prefix}.{device_id}.points"
|
|
|
|
# ------------------------------------------------------------------
|
|
def publish(self, samples: List[dict]) -> int:
|
|
"""上行一批样本;Kafka 发送成功的记录从 spool ack 删除。
|
|
|
|
Returns:
|
|
本轮成功上行条数。
|
|
"""
|
|
if self._degraded:
|
|
# 降级模式:样本已在 spool,等待 Kafka 恢复后重发
|
|
return 0
|
|
assert self._producer is not None
|
|
ok = 0
|
|
for s in samples:
|
|
topic = self._topic(s["device_id"])
|
|
record = {"device_id": s["device_id"], "point_id": s["point_id"],
|
|
"value": s["value"], "ts": s["ts"]}
|
|
try:
|
|
self._producer.produce(
|
|
topic,
|
|
key=s["point_id"].encode("utf-8"),
|
|
value=__import__("json").dumps(record, ensure_ascii=False).encode("utf-8"),
|
|
callback=self._on_delivery,
|
|
)
|
|
# 已进入 producer 缓冲即视为“已提交”,flush 时确认删除
|
|
self._sent += 1
|
|
ok += 1
|
|
except Exception as exc:
|
|
self._failed += 1
|
|
logger.warning("Kafka 发送失败(%s),样本保留在 spool", exc)
|
|
self._flush()
|
|
return ok
|
|
|
|
def _flush(self) -> None:
|
|
try:
|
|
if self._producer is not None:
|
|
self._producer.flush(timeout=5)
|
|
except Exception as exc:
|
|
logger.warning("Kafka flush 异常(%s)", exc)
|
|
|
|
def _on_delivery(self, err, msg) -> None: # pragma: no cover - 回调路径
|
|
"""投递确认回调:确认成功则 ack 删除 spool 记录(断点续传核心)。"""
|
|
if err is not None:
|
|
self._failed += 1
|
|
logger.warning("Kafka 投递失败: %s", err)
|
|
return
|
|
try:
|
|
point_id = msg.key().decode("utf-8") if msg.key() else ""
|
|
# 投递成功即按 point_id ack 最早一条未确认记录(幂等消费容忍轻微重发)
|
|
self.spool.ack({"device_id": None, "point_id": point_id, "value": None, "ts": None})
|
|
except Exception:
|
|
pass # ack 失败仅导致重发,幂等消费可容忍
|
|
|
|
# ------------------------------------------------------------------
|
|
@property
|
|
def stats(self) -> dict:
|
|
return {"sent": self._sent, "failed": self._failed, "degraded": self._degraded}
|
|
|
|
def close(self) -> None:
|
|
self._flush()
|
|
if self._producer is not None:
|
|
try:
|
|
self._producer.flush(timeout=5)
|
|
except Exception:
|
|
pass
|