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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 16:24:50