如何处理Apache Flink/AWS Kinesis Connector中的无效JSON?
解决PyFlink读取Kinesis时JSON格式错误导致任务中断的问题
问题描述
我有一个基于Apache Flink(PyFlink)的应用,通过官方Flink连接器和Table API从AWS Kinesis流读取数据。但Kinesis流中偶尔存在格式错误的JSON,会触发java.lang.RuntimeException: java.io.IOException: Failed to deserialize JSON异常,导致任务中断。尝试用Python的try/except捕获无效,因为异常来自预打包的Flink/Kinesis Connector jar。请问如何配置应用或连接器以忽略格式错误的条目并继续运行?
应用代码
from pyflink.table import EnvironmentSettings, TableEnvironment import os import json # 1. Creates a Table Environment env_settings = EnvironmentSettings.in_streaming_mode() table_env = TableEnvironment.create(env_settings) statement_set = table_env.create_statement_set() APPLICATION_PROPERTIES_FILE_PATH = "/etc/flink/application_properties.json" # on kda is_local = ( True if os.environ.get("IS_LOCAL") else False ) # set this env var in your local environment if is_local: # only for local, overwrite variable to properties and pass in your jars delimited by a semicolon (;) APPLICATION_PROPERTIES_FILE_PATH = "application_properties.json" # local CURRENT_DIR = os.path.dirname(os.path.realpath(__file__)) table_env.get_config().get_configuration().set_string( "pipeline.jars", "file:///" + CURRENT_DIR + "/lib/flink-sql-connector-kinesis-1.15.2.jar", ) def get_application_properties(): if os.path.isfile(APPLICATION_PROPERTIES_FILE_PATH): with open(APPLICATION_PROPERTIES_FILE_PATH, "r") as file: contents = file.read() properties = json.loads(contents) return properties else: print('A file at "{}" was not found'.format(APPLICATION_PROPERTIES_FILE_PATH)) def property_map(props, property_group_id): for prop in props: if prop["PropertyGroupId"] == property_group_id: return prop["PropertyMap"] def create_source_table(table_name, stream_name, region, stream_initpos): return """ CREATE TABLE {0} ( ticker VARCHAR(6), price DOUBLE, event_time TIMESTAMP(3), WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND ) PARTITIONED BY (ticker) WITH ( 'connector' = 'kinesis', 'stream' = '{1}', 'aws.region' = '{2}', 'scan.stream.initpos' = '{3}', 'format' = 'json', 'json.timestamp-format.standard' = 'ISO-8601' ) """.format( table_name, stream_name, region, stream_initpos ) def create_sink_table(table_name, stream_name, region, stream_initpos): return """ CREATE TABLE {0} ( ticker VARCHAR(6), price DOUBLE, event_time TIMESTAMP(3), WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND ) PARTITIONED BY (ticker) WITH ( 'connector' = 'kinesis', 'stream' = '{1}', 'aws.region' = '{2}', 'sink.partitioner-field-delimiter' = ';', 'sink.batch.max-size' = '100', 'format' = 'json', 'json.timestamp-format.standard' = 'ISO-8601' ) """.format( table_name, stream_name, region ) def create_print_table(table_name, stream_name, region, stream_initpos): return """ CREATE TABLE {0} ( ticker VARCHAR(6), price DOUBLE, event_time TIMESTAMP(3), WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND ) WITH ( 'connector' = 'print' ) """.format( table_name, stream_name, region, stream_initpos ) def main(): # Application Property Keys input_property_group_key = "consumer.config.0" producer_property_group_key = "producer.config.0" input_stream_key = "input.stream.name" input_region_key = "aws.region" input_starting_position_key = "flink.stream.initpos" output_stream_key = "output.stream.name" output_region_key = "aws.region" # tables input_table_name = "input_table" output_table_name = "output_table" # get application properties props = get_application_properties() input_property_map = property_map(props, input_property_group_key) output_property_map = property_map(props, producer_property_group_key) input_stream = input_property_map[input_stream_key] input_region = input_property_map[input_region_key] stream_initpos = input_property_map[input_starting_position_key] output_stream = output_property_map[output_stream_key] output_region = output_property_map[output_region_key] # 2. Creates a source table from a Kinesis Data Stream table_env.execute_sql( create_source_table(input_table_name, input_stream, input_region, stream_initpos) ) # 3. Creates a sink table writing to a Kinesis Data Stream table_env.execute_sql( create_sink_table(output_table_name, output_stream, output_region, stream_initpos) ) # 4. Inserts the source table data into the sink table table_result = table_env.execute_sql("INSERT INTO {0} SELECT * FROM {1}" .format(output_table_name, input_table_name)) if is_local: table_result.wait() else: # get job status through TableResult print(table_result.get_job_client().get_job_status()) if __name__ == "__main__": main()
解决方案
要忽略格式错误的JSON条目并让任务持续运行,只需在Kinesis源表的JSON格式配置中添加容错参数:
1. 添加JSON解析容错配置
修改create_source_table函数中的WITH配置块,加入'json.ignore-parse-errors' = 'true'。这个参数会让Flink跳过所有解析失败的JSON记录,继续处理后续有效数据。
修改后的create_source_table函数如下:
def create_source_table(table_name, stream_name, region, stream_initpos): return """ CREATE TABLE {0} ( ticker VARCHAR(6), price DOUBLE, event_time TIMESTAMP(3), WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND ) PARTITIONED BY (ticker) WITH ( 'connector' = 'kinesis', 'stream' = '{1}', 'aws.region' = '{2}', 'scan.stream.initpos' = '{3}', 'format' = 'json', 'json.timestamp-format.standard' = 'ISO-8601', 'json.ignore-parse-errors' = 'true' ) """.format( table_name, stream_name, region, stream_initpos )
2. 可选:补充字段缺失容错(按需)
如果流中还存在字段缺失的情况,可以额外添加'json.fail-on-missing-field' = 'false',避免因字段缺失触发异常:
'json.fail-on-missing-field' = 'false'
说明
通过上述配置,Flink会自动跳过解析失败的JSON条目,不会因为单条坏数据导致整个任务中断。如果需要对错误数据进行记录或后续分析,可以结合Flink的侧输出功能,将错误数据路由到单独的存储或流中。
内容的提问来源于stack exchange,提问作者m.clearpark
相关产品推荐
相关产品推荐

