Flink数据流结果能否存入List/Map?多流同步Hudi问题求助
基于Flink实现动态ID过滤的Hudi同步方案
问题根源说明
你之前的两种方案失败原因:
- 方案一:Flink的Source必须在作业初始化阶段(客户端构建DAG时)定义,不能在运行时的算子(如map)中动态创建,因此在map里新增的Source无法被Flink调度器识别,导致无数据输出。
- 方案二:Flink的作业执行流程是先构建DAG再提交运行,你写的for循环属于客户端构建DAG阶段的代码,而流读取ID并填充List是TaskManager运行时的操作,循环执行时List还未被填充,自然为空。
可行解决方案
方案1:使用广播状态(Broadcast State)
适合数据库ID动态更新的场景(如CDC实时同步ID变化),通过广播状态将ID集合分发到所有业务流Task,实现动态过滤。
步骤代码示例(Java):
import org.apache.flink.api.common.state.MapStateDescriptor; import org.apache.flink.streaming.api.datastream.BroadcastStream; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.streaming.api.functions.co.BroadcastProcessFunction; import org.apache.flink.streaming.connectors.kafka.KafkaSource; import org.apache.flink.connector.jdbc.JdbcSource; import org.apache.flink.util.Collector; // 初始化Flink环境 StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); // 1. 读取数据库ID的CDC流(实时获取ID新增/变更) DataStream<String> idStream = env.fromSource( JdbcSource.<String>builder() .setDrivername("com.mysql.cj.jdbc.Driver") .setUrl("jdbc:mysql://localhost:3306/your_db") .setUsername("root") .setPassword("your_pwd") .setQuery("SELECT id FROM your_id_table") .setRowTypeInfo(new org.apache.flink.api.java.typeutils.RowTypeInfo(org.apache.flink.api.common.typeinfo.BasicTypeInfo.STRING_TYPE_INFO)) .build(), org.apache.flink.api.common.eventtime.WatermarkStrategy.noWatermarks(), "JDBC-ID-Source" ).map(row -> row.getField(0).toString()); // 2. 定义广播状态描述符,存储ID集合 MapStateDescriptor<String, Boolean> idBroadcastStateDesc = new MapStateDescriptor<>( "id-broadcast-state", org.apache.flink.api.common.typeinfo.BasicTypeInfo.STRING_TYPE_INFO, org.apache.flink.api.common.typeinfo.BasicTypeInfo.BOOLEAN_TYPE_INFO ); // 3. 将ID流转为广播流 BroadcastStream<String> broadcastIdStream = idStream.broadcast(idBroadcastStateDesc); // 4. 读取需要过滤的业务数据流(示例为Kafka源) DataStream<BusinessData> businessStream = env.fromSource( KafkaSource.<BusinessData>builder() .setBootstrapServers("localhost:9092") .setTopics("business_topic") .setGroupId("hudi_sync_group") .setValueOnlyDeserializer(new org.apache.flink.api.common.serialization.JsonDeserializationSchema<>(BusinessData.class)) .build(), org.apache.flink.api.common.eventtime.WatermarkStrategy.forMonotonousTimestamps(), "Kafka-Business-Source" ); // 5. 连接业务流与广播流,实现动态过滤 DataStream<BusinessData> filteredStream = businessStream .connect(broadcastIdStream) .process(new BroadcastProcessFunction<BusinessData, String, BusinessData>() { @Override public void processElement(BusinessData value, ReadOnlyContext ctx, Collector<BusinessData> out) throws Exception { // 从广播状态获取ID集合,判断当前数据ID是否匹配 org.apache.flink.api.common.state.ReadOnlyBroadcastState<String, Boolean> broadcastState = ctx.getBroadcastState(idBroadcastStateDesc); if (broadcastState.contains(value.getId())) { out.collect(value); } } @Override public void processBroadcastElement(String id, Context ctx, Collector<BusinessData> out) throws Exception { // 更新广播状态,新增/更新ID org.apache.flink.api.common.state.BroadcastState<String, Boolean> broadcastState = ctx.getBroadcastState(idBroadcastStateDesc); broadcastState.put(id, true); } }); // 6. 将过滤后的数据流同步到Hudi filteredStream.sinkTo( org.apache.flink.table.dataformat.hudi.HuDiSink.<BusinessData>builder() .setRecordKeyField("id") .setTableName("hudi_business_table") .setBasePath("/path/to/hudi/storage") .build() ); env.execute("ID-Filter-Hudi-Sync-Job");
方案2:双流JOIN(适合ID全量静态场景)
如果数据库ID是一次性全量读取(如快照数据),可通过双流JOIN保留匹配的业务数据。
步骤代码示例(Java):
import org.apache.flink.api.java.tuple.Tuple2; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows; import org.apache.flink.streaming.api.windowing.time.Time; StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); // 1. 全量读取数据库ID流 DataStream<Tuple2<String, Integer>> idStream = env.fromSource( JdbcSource.<org.apache.flink.types.Row>builder() .setDrivername("com.mysql.cj.jdbc.Driver") .setUrl("jdbc:mysql://localhost:3306/your_db") .setUsername("root") .setPassword("your_pwd") .setQuery("SELECT id FROM your_id_table") .setRowTypeInfo(new org.apache.flink.api.java.typeutils.RowTypeInfo(org.apache.flink.api.common.typeinfo.BasicTypeInfo.STRING_TYPE_INFO)) .build(), org.apache.flink.api.common.eventtime.WatermarkStrategy.noWatermarks(), "JDBC-ID-Snapshot-Source" ).map(row -> Tuple2.of(row.getField(0).toString(), 1)); // 2. 业务流转为以ID为Key的流 DataStream<Tuple2<String, BusinessData>> businessKeyedStream = businessStream .map(data -> Tuple2.of(data.getId(), data)) .assignTimestampsAndWatermarks(org.apache.flink.api.common.eventtime.WatermarkStrategy.forMonotonousTimestamps()); // 3. 双流JOIN,只保留ID匹配的数据 DataStream<BusinessData> joinedStream = businessKeyedStream .join(idStream) .where(Tuple2::f0) .equalTo(Tuple2::f0) .window(TumblingEventTimeWindows.of(Time.hours(1))) // 根据业务数据延迟设置窗口 .apply((businessTuple, idTuple) -> businessTuple.f1); // 4. 同步到Hudi joinedStream.sinkTo( org.apache.flink.table.dataformat.hudi.HuDiSink.<BusinessData>builder() .setRecordKeyField("id") .setTableName("hudi_business_table") .setBasePath("/path/to/hudi/storage") .build() ); env.execute("ID-Join-Hudi-Sync-Job");
内容的提问来源于stack exchange,提问作者han liu
相关产品推荐
相关产品推荐

