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

使用Flink从两个数据源查找缺失记录的问题求助

问题描述

拥有两个数据源:S3存储桶(有界、按天分区)和PostgreSQL表,二者均以uuid作为唯一标识符,需求是找出S3中存在但PostgreSQL表中缺失的记录。尝试过程中遇到以下问题:

  • 使用Postgres Catalog读取PG时,因Flink 1.15.2不支持PG原生uuid类型报错
  • 改用JdbcInputFormat读取PG后,左连接查询因输出更新/删除变更导致Sink不兼容;改用toChangelogStream后输出包含+I(插入)、-D(删除)标识,不符合需求
  • 尝试批处理模式时,报错JdbcInputFormat创建的是无界表,不允许在批模式下查询

方案一:创建有界JDBC数据源(批处理模式)

Flink 1.15.2可以通过JDBC Connector明确配置有界数据源,结合批处理环境实现需求。

实现代码

// 初始化批处理环境
ExecutionEnvironment batchEnv = ExecutionEnvironment.getExecutionEnvironment();
BatchTableEnvironment tableEnv = BatchTableEnvironment.create(batchEnv);

// 读取S3数据(保持原有批处理逻辑)
final FileSource<GenericRecord> source = FileSource.forRecordStreamFormat(
        AvroParquetReaders.forGenericRecord(schema), path).build();
final DataStream<GenericRecord> avroStream = batchEnv.fromSource(
        source, WatermarkStrategy.noWatermarks(), "s3-source");
DataStream<Row> s3Stream = avroStream.map(x -> Row.of(x.get("uuid").toString()))
        .returns(Types.ROW_NAMED(new String[] {"uuid"}, Types.STRING));
Table s3table = tableEnv.fromDataStream(s3Stream);
tableEnv.createTemporaryView("s3table", s3table);

// 配置有界JDBC表源(通过DDL指定)
String createDbTableSql = "CREATE TEMPORARY TABLE dbtable (" +
        " uuid STRING" +
        ") WITH (" +
        " 'connector' = 'jdbc'," +
        " 'url' = 'jdbc:postgresql://127.0.0.1:5432/localdatabase'," +
        " 'table-name' = '(select cast(uuid as varchar) from localschema.table)'," +
        " 'username' = 'postgres'," +
        " 'password' = 'postgres'," +
        " 'driver' = 'org.postgresql.Driver'," +
        " 'scan.bounded.mode' = 'true'" + // 核心配置:标记为有界数据源
        ")";
tableEnv.executeSql(createDbTableSql);

// 执行左连接查询,获取缺失的uuid
Table resultTable = tableEnv.sqlQuery("SELECT s3table.uuid FROM s3table LEFT JOIN dbtable ON s3table.uuid = dbtable.uuid WHERE dbtable.uuid IS NULL");
DataStream<Row> resultStream = tableEnv.toDataStream(resultTable);
resultStream.print();

关键说明

  • 使用BatchTableEnvironment确保批处理语义,避免流模式的变更输出问题
  • 通过JDBC Connector的scan.bounded.mode参数明确将PG数据源标记为有界,解决批模式下的无界表报错
  • 直接在DDL的子查询中将PG的uuid转为字符串,绕过Flink对PG原生uuid类型的不支持问题

方案二:流API实现(过滤变更标识)

如果必须使用流模式,可以通过状态编程过滤掉变更标识,只保留最终的缺失uuid列表。

实现代码

// 初始化流环境
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);

// 读取S3数据(原有逻辑不变)
final FileSource<GenericRecord> source = FileSource.forRecordStreamFormat(
        AvroParquetReaders.forGenericRecord(schema), path).build();
final DataStream<GenericRecord> avroStream = env.fromSource(
        source, WatermarkStrategy.noWatermarks(), "s3-source");
DataStream<Row> s3Stream = avroStream.map(x -> Row.of(x.get("uuid").toString()))
        .returns(Types.ROW_NAMED(new String[] {"uuid"}, Types.STRING));
Table s3table = tableEnv.fromDataStream(s3Stream);
tableEnv.createTemporaryView("s3table", s3table);

// 读取PG数据(原有JdbcInputFormat逻辑不变)
TypeInformation<?>[] fieldTypes = new TypeInformation<?>[] {
        BasicTypeInfo.of(String.class)
};
RowTypeInfo rowTypeInfo = new RowTypeInfo(fieldTypes);
JdbcInputFormat jdbcInputFormat = JdbcInputFormat.buildJdbcInputFormat()
        .setDrivername("org.postgresql.Driver")
        .setDBUrl("jdbc:postgresql://127.0.0.1:5432/localdatabase")
        .setQuery("select cast(uuid as varchar) from localschema.table")
        .setUsername("postgres")
        .setPassword("postgres")
        .setRowTypeInfo(rowTypeInfo)
        .finish();
DataStream<Row> dbStream = env.createInput(jdbcInputFormat);
Table dbtable = tableEnv.fromDataStream(dbStream).as("uuid");
tableEnv.createTemporaryView("dbtable", dbtable);

// 执行左连接并转换为变更流
Table resultTable = tableEnv.sqlQuery("SELECT s3table.uuid FROM s3table LEFT JOIN dbtable ON s3table.uuid = dbtable.uuid WHERE dbtable.uuid IS NULL");
DataStream<RowData> changelogStream = tableEnv.toChangelogStream(resultTable);

// 过滤变更,只保留最终的缺失uuid
changelogStream
        .keyBy(row -> row.getString(0))
        .process(new KeyedProcessFunction<String, RowData, String>() {
            // 用状态记录已输出的uuid,避免重复或删除操作
            private final MapState<String, Boolean> seenUuids = getRuntimeContext().getMapState(
                    new MapStateDescriptor<>("seenUuids", String.class, Boolean.class));

            @Override
            public void processElement(RowData value, Context ctx, Collector<String> out) throws Exception {
                String uuid = value.getString(0);
                // 只处理首次出现的插入记录,忽略删除和重复项
                if (!seenUuids.contains(uuid)) {
                    out.collect(uuid);
                    seenUuids.put(uuid, true);
                }
            }
        })
        .print();

env.execute("Find Missing UUIDs");

关键说明

  • 使用toChangelogStream获取变更流后,通过KeyedProcessFunction维护状态
  • 利用MapState记录已经输出过的uuid,确保最终只输出唯一的缺失uuid,自动过滤-D标识和重复项

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 04:40:35