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
相关产品推荐
相关产品推荐

