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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 07:43:18