PyFlink中Kafka流关联静态数据时如何避免重复记录?
解决PyFlink中流与静态CSV关联的重复记录及时间序列化问题
一、重复记录问题解决
问题原因
流-批左关联会生成多条匹配中间结果,后续分组取最长前缀(max Code)的逻辑触发Flink状态更新;同时Flink流处理默认采用更新模式,当有延迟数据或状态变更时,会输出回撤消息(-U)和新插入消息(+I),导致终端呈现重复记录。
解决方案
1. 追加模式+ROW_NUMBER()严格去重
利用UniqueRecordID全局唯一性,通过ROW_NUMBER()过滤每个ID的最长匹配记录,并将输出切换为追加模式避免回撤消息:
WITH cte AS ( SELECT SetupTime, CallingNumber, CalledNumber, UniqueRecordID, Zone, CAST(Code AS DOUBLE) AS Code, Rate, sd.FileName AS FileName, ROW_NUMBER() OVER (PARTITION BY UniqueRecordID ORDER BY CAST(Code AS DOUBLE) DESC) AS rn FROM streaming_data sd LEFT JOIN static_data st ON sd.CalledNumber LIKE CONCAT(st.Code, '%') ) SELECT CalledNumber, Zone, CAST(Code AS VARCHAR(30)) AS Code, CAST(Rate AS VARCHAR(30)) AS Rate, SetupTime, UniqueRecordID, FileName FROM cte WHERE rn = 1
Table API中可通过配置开启追加模式:
ts_env.get_config().set("table.exec.sink.mode", "append")
2. 事件时间窗口内完成关联去重
针对有迟到数据的场景,基于SetupTime定义滚动窗口,在窗口内完成关联和最长匹配,确保每个窗口内的UniqueRecordID仅输出一次:
from pyflink.table import expressions as expr from pyflink.table.window import Tumble # 绑定事件时间窗口 streaming_window = ts_env.from_path("streaming_data").window( Tumble.over(expr.interval("5", "minutes")).on(expr.col("SetupTime")).alias("tumble_win") ) # 关联静态数据并标记窗口 cte = streaming_window.left_join(ts_env.from_path("static_data")) \ .where(expr.col("CalledNumber").like(expr.concat(expr.col("Code"), "%"))) \ .select( expr.col("SetupTime"), expr.col("CalledNumber"), expr.col("UniqueRecordID"), expr.col("Zone"), expr.col("Code").cast(expr.DOUBLE()).alias("Code"), expr.col("Rate"), expr.col("FileName"), expr.col("tumble_win") ) # 窗口内分组取最长匹配 windowed_max = cte.group_by(expr.col("tumble_win"), expr.col("UniqueRecordID")) \ .select( expr.col("UniqueRecordID"), expr.max(expr.col("Code")).alias("max_code"), expr.col("tumble_win") ) # 关联回静态数据获取完整字段 final_result = windowed_max.join(ts_env.from_path("static_data")) \ .where(expr.col("max_code") == expr.col("Code").cast(expr.DOUBLE())) \ .select( expr.col("CalledNumber"), expr.col("Zone"), expr.col("Code").cast(expr.STRING()).alias("Code"), expr.col("Rate").cast(expr.STRING()).alias("Rate"), expr.col("SetupTime"), expr.col("UniqueRecordID"), expr.col("FileName") ) final_result.execute().print()
3. UDF预加载静态数据优化匹配(推荐)
将静态CSV数据预加载到UDF中,直接在UDF内实现最长前缀匹配,避免流-批关联产生大量中间数据:
from pyflink.table import ScalarFunction, DataTypes from pyflink.table.udf import udf # 预加载静态数据到字典 static_dict = {} with open("static_data.csv", "r") as f: next(f) # 跳过表头 for line in f: zone, code, rate = line.strip().split(",") static_dict[code] = (zone, rate) # 自定义最长前缀匹配UDF class LongestMatchUDF(ScalarFunction): def eval(self, called_number): max_len = -1 result = (None, None) for code in static_dict.keys(): if called_number.startswith(code) and len(code) > max_len: max_len = len(code) result = static_dict[code] return result # 注册UDF longest_match_udf = udf(LongestMatchUDF(), result_type=DataTypes.ROW([ DataTypes.FIELD("Zone", DataTypes.STRING()), DataTypes.FIELD("Rate", DataTypes.STRING()) ])) # 直接调用UDF处理流数据 final_result = ts_env.from_path("streaming_data") \ .select( expr.col("CalledNumber"), longest_match_udf(expr.col("CalledNumber")).alias("match_res"), expr.col("SetupTime"), expr.col("UniqueRecordID"), expr.col("FileName") ) \ .select( expr.col("CalledNumber"), expr.col("match_res.Zone").alias("Zone"), expr.col("match_res.Rate").alias("Rate"), expr.col("SetupTime"), expr.col("UniqueRecordID"), expr.col("FileName") ) final_result.execute().print()
二、时间字段序列化错误解决
错误1:DateTimeParseException(格式不匹配)
原因:Flink默认TIMESTAMP解析格式为yyyy-MM-dd'T'HH:mm:ss.SSS,但你的数据格式是yyyy-MM-dd HH:mm:ss.SSS(空格分隔),导致解析失败。
解决方案:定义streaming_data表时指定时间解析格式:
CREATE TABLE streaming_data ( SetupTime TIMESTAMP(3), CallingNumber STRING, CalledNumber STRING, UniqueRecordID STRING, FileName STRING ) WITH ( 'connector' = 'kafka', 'topic' = 'your_topic', 'properties.bootstrap.servers' = 'localhost:9092', 'format' = 'json', 'json.timestamp-format.standard' = 'yyyy-MM-dd HH:mm:ss.SSS' );
错误2:ClassCastException(类型不兼容)
原因:Types.SQL_TIMESTAMP()对应java.sql.Timestamp,而TIMESTAMP_LTZ对应java.time.Instant,类型映射不匹配导致转换错误。
解决方案:
- 本地时间(无时区)使用
DataTypes.TIMESTAMP(3); - UTC时间使用
DataTypes.TIMESTAMP_LTZ(3); - 统一使用
DataTypes定义表结构,避免混合Types与DataTypes:
from pyflink.table import DataTypes ts_env.create_temporary_table("streaming_data", DataTypes.ROW([ DataTypes.FIELD("SetupTime", DataTypes.TIMESTAMP(3)), DataTypes.FIELD("CallingNumber", DataTypes.STRING()), DataTypes.FIELD("CalledNumber", DataTypes.STRING()), DataTypes.FIELD("UniqueRecordID", DataTypes.STRING()), DataTypes.FIELD("FileName", DataTypes.STRING()) ]), { 'connector': 'kafka', 'topic': 'your_topic', 'properties.bootstrap.servers': 'localhost:9092', 'format': 'json', 'json.timestamp-format.standard': 'yyyy-MM-dd HH:mm:ss.SSS' })
内容的提问来源于stack exchange,提问作者user8045747
相关产品推荐
相关产品推荐

