使用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
相关产品推荐
相关产品推荐

