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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 01:22:05