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

PyFlink在Kinesis Analytics Studio中无法转换DataStream至Kinesis流

问题描述

我有一个基于CoFlatMapFunction生成的DataStream <pyflink.datastream.data_stream.DataStream>,代码简化如下:

%flink.pyflink
# 连接两个流并更新规则集
class MyCoFlatMapFunction(CoFlatMapFunction):

    def open(self, runtime_context: RuntimeContext):
        state_desc = MapStateDescriptor('map', Types.STRING(), Types.BOOLEAN())
        self.state = runtime_context.get_map_state(state_desc)

    def bool_from_user_number(self, user_number: int):
        '''当user_number大于0时返回True,否则返回False'''
        if user_number > 0:
            return True
        else:
            return False

    def flat_map1(self, value):
        '''处理第一个连接流中的每个元素'''
        self.state.put(value[1], self.bool_from_user_number(value[2]))

    def flat_map2(self, value):
        '''处理第二个连接流(exchange_server_tickers_data_py)中的每个元素'''
        current_dateTime = datetime.now()
        dt = current_dateTime

        x = value[1]
        y = value[2]

        yield Row(dt, x, y)

def generate__ds(st_env):
    # 将动态更新的表转换为DataStream
    type_info1 = Types.ROW([Types.SQL_TIMESTAMP(), Types.STRING(), Types.INT()])
    ds1 = st_env.to_append_stream(table_1 , type_info=type_info1)

    type_info2 = Types.ROW([Types.SQL_TIMESTAMP(), Types.STRING(), Types.STRING()])
    ds2 = st_env.to_append_stream(table_2 , type_info=type_info2)

    output_type_info = Types.ROW([ Types.PICKLED_BYTE_ARRAY() ,Types.STRING(),Types.STRING() ])
    # 连接两个流
    connected_ds = ds1.connect(ds2)
    # 应用CoFlatMapFunction
    ds = connected_ds.key_by(lambda a: a[0], lambda a: a[0]).flat_map(MyCoFlatMapFunction(), output_type_info)
    return ds

ds = generate__ds(st_env)

但我无法查看输出结果,不管是将其注册为视图/表、写入sink表,还是(最优方案)用Kinesis Streams sink把Flink流的数据写入Kinesis流都不行。Firehose的30秒延迟太长,不符合我的需求。

已尝试的操作

  1. 注册为视图/表
# 将DataStream转换为Table
input_table = st_env.from_data_stream(ds).alias("dt", "x", "y")
z.show(input_table, stream_type="update")

报错:

Query schema: [dt: RAW('[B', '...'), x: STRING, y: STRING]
Sink schema: [dt: RAW('[B', ?), x: STRING, y: STRING]

  1. 写入sink表
%flink.pyflink
# 创建用于输出结果的sink表
st_env.execute_sql("""DROP TABLE IF EXISTS table_sink""")

st_env.execute_sql("""
    CREATE TABLE table_sink (
        dt RAW('[B', '...'),
        x VARCHAR(32),
        y STRING
    ) WITH (
        'connector' = 'print'
    )
""")

# 将Table API的表转换为SQL视图
table = st_env.from_data_stream(ds).alias("dt", "spread", "spread_orderbook")
st_env.execute_sql("""DROP TEMPORARY VIEW IF EXISTS table_api_table""")
st_env.create_temporary_view('table_api_table', table)

# 输出Table API表的数据
st_env.execute_sql("INSERT INTO table_sink SELECT * FROM table_api_table").wait()

报错:

org.apache.flink.table.api.ValidationException: Unable to restore the RAW type of class '[B' with serializer snapshot '...'.

  1. 使用sink_to写入Kinesis流
%flink.pyflink
from pyflink.common.serialization import JsonRowSerializationSchema
from pyflink.datastream.connectors import KinesisStreamsSink

output_type_info = Types.ROW([Types.SQL_TIMESTAMP(), Types.STRING(), Types.STRING()])
serialization_schema = JsonRowSerializationSchema.Builder().with_type_info(output_type_info).build()

# 必填配置
sink_properties = {
    'aws.region': 'eu-west-2'
}

kds_sink = KinesisStreamsSink.builder()
            .set_kinesis_client_properties(sink_properties)
            .set_serialization_schema(SimpleStringSchema())
            .set_partition_key_generator(PartitionKeyGenerator
            .fixed())
            .set_stream_name("test_stream")
            .set_fail_on_error(False)
            .set_max_batch_size(500)
            .set_max_in_flight_requests(50)
            .set_max_buffered_requests(10000)
            .set_max_batch_size_in_bytes(5 * 1024 * 1024)
            .set_max_time_in_buffer_ms(5000)
            .set_max_record_size_in_bytes(1 * 1024 * 1024)
            .build()

ds.sink_to(kds_sink)

但在pyflink.datastream.connectors中找不到KinesisStreamsSink,也找不到AWS Kinesis Analytics Studio中的相关文档。

请问如何将数据写入Kinesis Streams sink,或者如何将该DataStream正确转换为表?


解决方案

一、修复DataStream转Table的问题

问题根源在于你定义的output_type_info使用了Types.PICKLED_BYTE_ARRAY()存储datetime类型,导致Flink无法正确序列化/反序列化这个RAW类型。需要修改输出类型为Flink支持的标准时间类型:

  1. 调整CoFlatMapFunction的输出类型
    在generate__ds函数中,把output_type_info改成:

    output_type_info = Types.ROW([Types.SQL_TIMESTAMP(), Types.STRING(), Types.STRING()])
    

    同时确保flat_map2中生成的dt是符合SQL_TIMESTAMP格式的类型,建议转换为Flink的Timestamp类型,避免本地datetime对象的序列化问题:

    from pyflink.common import types
    # 在flat_map2中替换dt的生成
    dt = types.Timestamp.from_datetime(datetime.now().astimezone())
    
  2. 重新转换为Table并查看输出
    调整后再注册视图或写入sink表就不会有RAW类型的报错了:

    input_table = st_env.from_data_stream(ds).alias("dt", "x", "y")
    z.show(input_table, stream_type="append")  # 用append模式更适配你的场景
    

    或者写入print sink:

    st_env.execute_sql("""
        CREATE TABLE table_sink (
            dt TIMESTAMP(3),
            x VARCHAR(32),
            y STRING
        ) WITH (
            'connector' = 'print'
        )
    """)
    table = st_env.from_data_stream(ds).alias("dt", "x", "y")
    st_env.create_temporary_view('table_api_table', table)
    st_env.execute_sql("INSERT INTO table_sink SELECT * FROM table_api_table").wait()
    

二、实现Kinesis Streams Sink的正确方式

在AWS Kinesis Analytics Studio中,推荐使用Table API/SQL的Kinesis连接器而非DataStream API的sink,托管环境中前者配置更简便:

方式1:使用Table SQL创建Kinesis Sink表

CREATE TABLE kinesis_sink (
    dt TIMESTAMP(3),
    x VARCHAR(32),
    y STRING
) WITH (
    'connector' = 'kinesis',
    'stream' = 'test_stream',
    'aws.region' = 'eu-west-2',
    'format' = 'json',
    'sink.partitioner' = 'fixed' -- 根据业务需求选择分区策略
);

然后将DataStream转换为Table后写入:

table = st_env.from_data_stream(ds).alias("dt", "x", "y")
st_env.execute_sql("INSERT INTO kinesis_sink SELECT * FROM table").wait()

方式2:使用DataStream API的Kinesis Sink

注意PyFlink中Kinesis Streams Sink的正确导入路径是pyflink.datastream.connectors.kinesis,而非直接在connectors下:

from pyflink.datastream.connectors.kinesis import KinesisStreamsSink, KinesisSinkConfiguration
from pyflink.common.serialization import JsonRowSerializationSchema
from pyflink.common.typeinfo import Types

# 确保DataStream的输出类型是标准类型
output_type_info = Types.ROW([Types.SQL_TIMESTAMP(), Types.STRING(), Types.STRING()])

# 配置Kinesis Sink
sink_config = KinesisSinkConfiguration()
sink_config.set_property('aws.region', 'eu-west-2')

# 序列化Schema
serialization_schema = JsonRowSerializationSchema.Builder()
    .with_type_info(output_type_info)
    .build()

kds_sink = KinesisStreamsSink.builder()
    .set_stream_name("test_stream")
    .set_kinesis_sink_configuration(sink_config)
    .set_serialization_schema(serialization_schema)
    .set_partition_key_generator(lambda x: x[1])  # 自定义分区键逻辑,也可使用fixed()
    .build()

# 添加sink并执行
ds.sink_to(kds_sink)
st_env.execute()

托管环境通常默认包含Kinesis连接器依赖,自定义环境可能需要手动添加依赖包。


内容的提问来源于stack exchange,提问作者Shav Vimalendiran

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 22:40:44