kafka-python库中Kafka记录message.timestamp赋值为哪个时间值?
kafka-python 消费记录
.timestamp 属性取值规则 kafka-python本身没有自定义这个字段的取值逻辑,完全对齐Apache Kafka原生协议的消息时间戳规范,具体取值和消费目标topic的配置直接相关:
- 默认场景下,topic的
message.timestamp.type参数配置为CreateTime,此时.timestamp存储的是生产者客户端构造消息时,在本地打上的消息创建时间,为毫秒级Unix时间戳。这个值依赖生产者本地时钟,如果生产者节点时钟偏移,拿到的时间值可能和broker、消费者侧的当前时间存在偏差。 - 如果topic的
message.timestamp.type参数被修改为LogAppendTime,此时.timestamp存储的是消息被broker接收并成功写入对应分区日志时,broker节点打上的日志追加时间,同样为毫秒级Unix时间戳。这种模式下同分区内消息的时间戳严格单调递增,不会出现后写入的消息时间戳更小的情况,时间值以broker节点时钟为准。
你可以同时通过记录对象的.timestamp_type属性判断当前时间戳的类型:
- 返回值为
0时,代表当前时间戳是CreateTime - 返回值为
1时,代表当前时间戳是LogAppendTime
补充几个容易踩的细节:
- 压缩消息场景下,CreateTime模式中批次外层的压缩包时间戳会取批次内所有消息的最大创建时间,kafka-python在解压还原单条消息时,会自动把单条消息原本的时间戳赋值给
.timestamp属性,不需要额外手动解析。 - 事务消息场景下,不管事务提交延迟多久,消息的
.timestamp取值仍然遵循上述两种规则,不会被替换为消费者读取到消息的时间,也不会被替换为事务提交的时间。 - 注意:这个字段永远不会赋值为消费者客户端本地收到消息的时间,如果需要统计消费端处理延迟,需要在消费逻辑执行时自行采集本地时间戳做差值计算。
内容的提问来源于stack exchange,提问作者mirakuru970
相关产品推荐
相关产品推荐

