使用Flink Python API将含时间戳数据写入GCS Parquet遇类型转换问题
问题
我有一个存储JSON消息的Kafka Topic,尝试通过Flink Python API处理这些消息并将其存储到GCS的Parquet文件中。以下是整理后的代码片段:
class Extract(MapFunction): def map(self, value): record = json.loads(value) dt_object = datetime.strptime(record['ts'], "%Y-%m-%dT%H:%M:%SZ") return Row(dt_object, record['event_id']) <...> events_schema = DataTypes.ROW([ DataTypes.FIELD("ts", DataTypes.TIMESTAMP()), DataTypes.FIELD("event_id", DataTypes.STRING()) ]) <...> # Main job part kafka_source = KafkaSource.builder() \ <...> .build() ds: DataStream = env.from_source(kafka_source, WatermarkStrategy.no_watermarks(), "Kafka Source") mapped_data = ds.map(Extract(), Types.ROW([Types.SQL_TIMESTAMP(), Types.STRING()])) sink = (FileSink .for_bulk_format("gs://<my_events_path>", ParquetBulkWriters.for_row_type(row_type=events_schema)) .with_output_file_config( OutputFileConfig.builder() .with_part_prefix("bids") .with_part_suffix(".parquet") .build()) .build()) mapped_data.sink_to(sink)
运行作业时出现异常:Java.lang.ClassCastException: class java.sql.Timestamp cannot be cast to class java.time.LocalDateTime (java.sql.Timestamp is in module java.sql of loader 'platform'; java.time.LocalDateTime is in module java.base of loader 'bootstrap'),原因是Types.SQL_TIMESTAMP()与DataTypes.TIMESTAMP()对应的Java类不兼容。若改为mapped_data = ds.map(Extract()),则会出现另一个异常:java.lang.ClassCastException: class [B cannot be cast to class org.apache.flink.types.Row ([B is in module java.base of loader 'bootstrap'; org.apache.flink.types.Row is in unnamed module of loader 'app')。想请教是否可以通过Flink Python API存储包含时间戳的Parquet格式数据?
解决方案
可以通过Flink Python API存储带时间戳的Parquet数据,问题出在时间戳类型匹配和Kafka Source反序列化配置上,修正步骤如下:
1. 统一Row对象的字段映射
原代码中返回的Row(dt_object, record['event_id'])是无名字段,需要和events_schema中的字段名一一对应,同时确保时间戳类型兼容:
class Extract(MapFunction): def map(self, value): record = json.loads(value) dt = datetime.strptime(record['ts'], "%Y-%m-%dT%H:%M:%SZ") # 显式指定字段名,与schema中的ts、event_id对应 return Row(ts=dt, event_id=record['event_id'])
2. 配置Kafka Source的字符串反序列化器
第二个异常是因为Kafka Source默认返回字节数组,需要添加SimpleStringSchema将字节转为字符串,才能被json.loads解析:
from pyflink.common.serialization import SimpleStringSchema kafka_source = KafkaSource.builder() \ .set_bootstrap_servers("<your-bootstrap-servers>") \ .set_topics("<your-topic>") \ .set_group_id("<your-group-id>") \ .set_starting_offsets(KafkaOffsetsInitializer.earliest()) \ # 添加字符串反序列化器,解决[B转Row的异常 .set_value_only_deserializer(SimpleStringSchema()) \ .build()
3. 匹配Map算子的输出类型与Sink Schema
使用DataTypes定义的schema作为map算子的输出类型,替代Types类,确保类型完全对齐:
mapped_data = ds.map(Extract(), output_type=events_schema)
完整修正后的核心代码
from datetime import datetime import json from pyflink.common import Row, DataTypes from pyflink.datastream import MapFunction, StreamExecutionEnvironment from pyflink.datastream.connectors.kafka import KafkaSource, KafkaOffsetsInitializer from pyflink.datastream.connectors.file_system import FileSink, OutputFileConfig from pyflink.common.serialization import SimpleStringSchema from pyflink.datastream.connectors.file_system.bulk import ParquetBulkWriters env = StreamExecutionEnvironment.get_execution_environment() class Extract(MapFunction): def map(self, value): record = json.loads(value) dt = datetime.strptime(record['ts'], "%Y-%m-%dT%H:%M:%SZ") return Row(ts=dt, event_id=record['event_id']) events_schema = DataTypes.ROW([ DataTypes.FIELD("ts", DataTypes.TIMESTAMP()), DataTypes.FIELD("event_id", DataTypes.STRING()) ]) kafka_source = KafkaSource.builder() \ .set_bootstrap_servers("your-bootstrap-servers") \ .set_topics("your-topic") \ .set_group_id("your-group-id") \ .set_starting_offsets(KafkaOffsetsInitializer.earliest()) \ .set_value_only_deserializer(SimpleStringSchema()) \ .build() ds = env.from_source(kafka_source, WatermarkStrategy.no_watermarks(), "Kafka Source") mapped_data = ds.map(Extract(), output_type=events_schema) sink = (FileSink .for_bulk_format("gs://<my_events_path>", ParquetBulkWriters.for_row_type(row_type=events_schema)) .with_output_file_config( OutputFileConfig.builder() .with_part_prefix("bids") .with_part_suffix(".parquet") .build()) .build()) mapped_data.sink_to(sink) env.execute("Kafka to GCS Parquet Job")
关键注意点
Row对象必须显式指定字段名,与events_schema中的字段完全匹配,Flink依赖字段名进行类型映射DataTypes.TIMESTAMP()可直接接收Python的datetime.datetime对象,无需手动转换为Java类型- 必须配置
SimpleStringSchema,否则Kafka返回的字节数组无法被JSON解析
内容的提问来源于stack exchange,提问作者xneg

