如何在PyFlink DataStream API中获取Kafka消息的Rowtime?
在PyFlink DataStream API中获取Kafka消息的AppendLogTime
要在DataStream API中获取Kafka消息的AppendLogTime(即Kafka Broker记录的消息写入时间),你之前的实现思路存在两个核心问题:
- 你使用的
JsonRowDeserializationSchema仅会解析Kafka消息体中的JSON数据,不会自动提取Kafka的元数据(包括AppendLogTime),所以value[5]只是消息体里的普通Long值,不存在rowtime属性,自然会触发AttributeError。 - TableAPI中的
.rowtime语法是用来将某个字段标记为事件时间属性,用于后续的时间窗口等操作,并非直接从Kafka元数据中提取时间戳。
正确实现方式
要获取Kafka的元数据,需要使用KafkaRecordDeserializationSchema,它可以同时获取消息体和Kafka的元数据(包括timestamp)。以下是完整的代码示例:
1. 定义自定义的KafkaRecordDeserializationSchema
from pyflink.common import Row from pyflink.common.typeinfo import Types from pyflink.datastream.connectors.kafka import KafkaRecordDeserializationSchema from pyflink.datastream.connectors.kafka import FlinkKafkaConsumer from pyflink.datastream import StreamExecutionEnvironment from pyflink.datastream.functions import MapFunction import json from datetime import datetime class CustomKafkaSchema(KafkaRecordDeserializationSchema): def __init__(self): # 定义消息体的字段和类型 self.row_type = Types.ROW_NAMED( ['name', 'platform', 'year', 'global_sales', 'time_send', 'time_in_sps', 'write_time'], [Types.STRING(), Types.STRING(), Types.INT(), Types.DOUBLE(), Types.STRING(), Types.STRING(), Types.STRING()] ) def deserialize(self, record): # 解析消息体的JSON数据 message_body = json.loads(record.value()) # 获取Kafka的AppendLogTime(record.timestamp()返回的就是Broker的写入时间) append_log_time = record.timestamp() # 将消息体字段和元数据组合成新的Row return Row( message_body['name'], message_body['platform'], message_body['year'], message_body['global_sales'], message_body['time_send'], append_log_time, # 这里就是Kafka的AppendLogTime message_body['time_in_sps'], message_body['write_time'] ) def get_produced_type(self): # 返回最终Row的类型信息,包含append_log_time字段 return Types.ROW_NAMED( ['name', 'platform', 'year', 'global_sales', 'time_send', 'append_log_time', 'time_in_sps', 'write_time'], [Types.STRING(), Types.STRING(), Types.INT(), Types.DOUBLE(), Types.STRING(), Types.LONG(), Types.STRING(), Types.STRING()] )
2. 创建Kafka Source并使用自定义Schema
env = StreamExecutionEnvironment.get_execution_environment() env.set_parallelism(1) kafka_props = { 'bootstrap.servers': 'kafka:9092', 'group.id': 'test_group' } kafka_source = FlinkKafkaConsumer( topics='test', deserialization_schema=CustomKafkaSchema(), properties=kafka_props ) ds = env.add_source(kafka_source)
3. 在MapFunction中直接访问AppendLogTime
此时append_log_time已经是Row中的一个普通Long字段,直接通过索引或字段名访问即可:
class MyMapFunction(MapFunction): def map(self, value): # 直接获取append_log_time,无需访问rowtime属性 current_time = str(datetime.timestamp(datetime.now()) * 1000) return Row( value[0], value[1], value[2], value[3], value[4], value[5], # 这里就是Kafka的AppendLogTime current_time, value[7] ) ds = ds.map(MyMapFunction()) ds.print() env.execute("Kafka AppendLogTime Demo")
关键说明
record.timestamp()返回的就是Kafka Broker记录的消息写入时间(AppendLogTime),这个值由Kafka自动生成,不需要消息体中包含该字段。- 如果需要区分不同类型的时间戳(比如消息创建时间、Broker写入时间),可以通过
record.timestamp_type()获取时间戳类型,判断是否为TimestampType.LOG_APPEND_TIME。
内容的提问来源于stack exchange,提问作者Hansanho
相关产品推荐
相关产品推荐

