You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

confluent-kafka[avro]中creationTime字段序列化后时间异常求助

问题:Avro时间字段序列化异常,显示为1970年早期时间

现象

使用confluent-kafka[avro]==2.1.1库向Kafka发送Avro格式消息时,字段creationTime始终显示为类似1970-01-20T13:34:12.378Z的错误时间:

  • 消息默认时间戳正常,但将creationTime设为消息时间戳时,消息显示为58年前发送
  • 所有环境(本地、AWS ECS)均可复现
  • 调试发现:序列化前timestamp是浮点型秒级时间戳(如1690451888.45323),但AKHQ接收后解析出的时间对应的时间戳仅为1686851(毫秒级)

相关代码与配置

获取时间戳的代码

timestamp = datetime.now(tz=tz.gettz(self.config.timezone)).timestamp()
values = {
    "creationTime": timestamp,
    "description": message,
    "eventTypeId": self.config.metric_name,
    "pgd": [],
    "contracts": [],
    "points": [],
    "objects": [],
}

Kafka生产者代码

"""This module contains everything necessary to send messages to kafka"""
import logging

from confluent_kafka import SerializingProducer
from confluent_kafka.schema_registry import SchemaRegistryClient
from confluent_kafka.schema_registry.avro import AvroSerializer
from confluent_kafka.serialization import StringSerializer

from param_store_models import KafkaInfo

LOGGER = logging.getLogger(__name__)


class KafkaProducer:
    """Class used to send messages to kafka"""

    def __init__(self, schema: str, kafka_info: KafkaInfo):
        producer_ssm_conf = {
            "bootstrap.servers": kafka_info.bootstrap_servers,
            "security.protocol": kafka_info.security_protocol,
            "sasl.mechanism": kafka_info.sasl_mecanism,
            "sasl.username": kafka_info.sasl_username,
            "sasl.password": kafka_info.sasl_password,
        }
        registry_ssm_conf = {"url": kafka_info.schema_registry_url}

        serializer = AvroSerializer(
            SchemaRegistryClient(registry_ssm_conf), schema, conf={"auto.register.schemas": False}
        )

        producer_default_conf = {
            "value.serializer": serializer,
            "key.serializer": StringSerializer(),
            "enable.idempotence": "true",
            "max.in.flight.requests.per.connection": 1,
            "retries": 5,
            "acks": "all",
            "retry.backoff.ms": 500,
            "queue.buffering.max.ms": 500,
            "error_cb": self.delivery_report,
        }

        self.__serializing_producer = SerializingProducer({**producer_default_conf, **producer_ssm_conf})

    def produce(self, topic: str, key=None, value=None, timestamp=0):
        """Asynchronously produce message to a topic"""
        LOGGER.info(f"Produce message {value} to topic {topic}")
        self.__serializing_producer.produce(topic, key, value, on_delivery=self.delivery_report, timestamp=timestamp)

    def flush(self):
        """
        Flush messages and trigger callbacks
        :return: Number of messages still in queue.
        """
        LOGGER.debug("Flushing messages to kafka")
        return self.__serializing_producer.flush()

    @staticmethod
    def delivery_report(err, msg):
        """
        Called once for each message produced to indicate delivery result.
        Triggered by poll() or flush().
        """
        if err:
            LOGGER.error(f"Kafka message delivery failed: {err}")
        else:
            LOGGER.info(f"Kafka message delivered to {msg.topic()} [{msg.partition()}]")

Avro Schema(脱敏后)

{
  "type": "record",
  "name": "EventRecord",
  "namespace": "com.event",
  "doc": "Schéma d'un évènement de supervision brut",
  "fields": [
    {
      "name": "creationTime",
      "type": {
        "type": "long",
        "logicalType": "timestamp-millis"
      }
    },
    {
      "name": "eventTypeId",
      "type": "string"
    },
    {
      "name": "internalProductId",
      "type": [
        "null",
        "string"
      ],
      "default": null
    },
    {
      "name": "description",
      "type": [
        "null",
        "string"
      ],
      "default": null
    },
    {
      "name": "contracts",
      "type": {
        "type": "array",
        "items": "string"
      }
    },
    {
      "name": "points",
      "type": {
        "type": "array",
        "items": "string"
      }
    },
    {
      "name": "objects",
      "type": {
        "type": "array",
        "items": "string"
      }
    }
  ]
}

AKHQ接收的消息示例

{
  "creationTime": "1970-01-20T13:34:12.378Z",
  "eventTypeId": "test",
  "description": "Test",
  "pgd": [],
  "contracts": [],
  "points": [],
  "objects": []
}

原因与解决方案

核心原因

Avro Schema中creationTime的logicalType: timestamp-millis要求传入从Epoch开始的毫秒级时间戳(整数类型),但当前代码中datetime.timestamp()返回的是秒级时间戳(浮点类型)。当浮点数被强制转换为long类型时,仅保留了整数部分(秒数),AKHQ会将这个值当作毫秒数解析,导致时间被缩小为原有的1/1000,从而显示为1970年的早期时间。

解决方案

修改时间戳生成代码,将秒级浮点数转换为毫秒级整数:

# 替换原timestamp生成代码
dt = datetime.now(tz=tz.gettz(self.config.timezone))
# 将秒级时间戳乘以1000后转为整数,得到毫秒级时间戳
timestamp = int(dt.timestamp() * 1000)

values = {
    "creationTime": timestamp,
    # 其他字段保持不变
}

更稳妥的精度处理方式

如果担心浮点精度问题,可以直接计算UTC时间的毫秒数:

from datetime import timezone

dt = datetime.now(tz=tz.gettz(self.config.timezone)).astimezone(timezone.utc)
timestamp = int(dt.timestamp() * 1000)

修改后,传入的creationTime将符合Avrotimestamp-millis的要求,AKHQ会解析出正确的当前时间。


内容的提问来源于stack exchange,提问作者Junn Sorran

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.14 20:15:56