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秒延迟太长,不符合我的需求。
已尝试的操作
- 注册为视图/表
# 将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]
- 写入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 '...'.
- 使用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支持的标准时间类型:
调整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())重新转换为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

