# -*- 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