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

如何在Flink中连接两个Kafka Streams?利用静态主表流丰富实时流并实现内存永久Lookup Table

嘿,这个场景我刚好实操过!针对你想用静态Kafka流作为内存Lookup Table来持续丰富实时流的需求,Flink里有两个非常贴合的方案,我给你详细拆解:

一、DataStream API:Broadcast State 方案

如果你的静态流是少量数据,并且希望每个Task节点都持有一份完整的Lookup数据(永久驻留内存,支持后续动态更新静态数据),Broadcast State是最优选择。它能把静态流的所有数据广播到Flink集群的每个Task,实时流到来时直接在本地内存做关联,性能拉满。

代码示例

import org.apache.flink.api.common.state.BroadcastState;
import org.apache.flink.api.common.state.MapStateDescriptor;
import org.apache.flink.api.common.state.ReadOnlyBroadcastState;
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.FlinkKafkaConsumer;
import org.apache.flink.util.Collector;

import java.util.Properties;

public class StaticLookupEnrichment {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        // 1. 配置Kafka消费者属性
        Properties kafkaProps = new Properties();
        kafkaProps.setProperty("bootstrap.servers", "localhost:9092");
        kafkaProps.setProperty("group.id", "flink-lookup-group");

        // 2. 读取实时数据流(比如订单流,key是用户ID)
        DataStream<Order> realTimeOrderStream = env.addSource(
                new FlinkKafkaConsumer<>("real-time-orders", new OrderDeserializationSchema(), kafkaProps)
        );

        // 3. 读取静态Lookup流(比如用户信息表,key是用户ID)
        DataStream<UserInfo> staticUserStream = env.addSource(
                new FlinkKafkaConsumer<>("static-user-info", new UserInfoDeserializationSchema(), kafkaProps)
        );

        // 4. 定义Broadcast State的描述符,用来存储用户信息(key: 用户ID, value: 用户详情)
        MapStateDescriptor<String, UserInfo> userLookupDescriptor =
                new MapStateDescriptor<>("user-lookup-table", String.class, UserInfo.class);

        // 5. 将静态流广播到所有Task
        BroadcastStream<UserInfo> broadcastUserStream = staticUserStream.broadcast(userLookupDescriptor);

        // 6. 关联实时流和广播的Lookup表
        DataStream<EnrichedOrder> enrichedOrderStream = realTimeOrderStream
                .connect(broadcastUserStream)
                .process(new BroadcastProcessFunction<Order, UserInfo, EnrichedOrder>() {
                    // 处理实时订单流
                    @Override
                    public void processElement(Order order, ReadOnlyContext ctx, Collector<EnrichedOrder> out) throws Exception {
                        // 从广播状态中获取对应用户的信息
                        ReadOnlyBroadcastState<String, UserInfo> lookupState = ctx.getBroadcastState(userLookupDescriptor);
                        UserInfo userInfo = lookupState.get(order.getUserId());

                        // 如果找到用户信息,就输出富化后的订单
                        if (userInfo != null) {
                            out.collect(new EnrichedOrder(
                                    order.getOrderId(),
                                    order.getUserId(),
                                    order.getAmount(),
                                    userInfo.getUserName(),
                                    userInfo.getUserLevel()
                            ));
                        }
                    }

                    // 处理静态流的更新(如果静态流后续有数据更新,会自动更新广播状态)
                    @Override
                    public void processBroadcastElement(UserInfo userInfo, Context ctx, Collector<EnrichedOrder> out) throws Exception {
                        BroadcastState<String, UserInfo> broadcastState = ctx.getBroadcastState(userLookupDescriptor);
                        // 用新的用户信息更新Lookup表(如果存在则覆盖,不存在则新增)
                        broadcastState.put(userInfo.getUserId(), userInfo);
                    }
                });

        // 输出富化后的结果(可以写到Kafka、数据库等)
        enrichedOrderStream.print();

        env.execute("Static Lookup Enrichment Job");
    }

    // 自定义实体类:订单
    public static class Order {
        private String orderId;
        private String userId;
        private double amount;
        // 省略getter、setter、构造方法
    }

    // 自定义实体类:用户信息
    public static class UserInfo {
        private String userId;
        private String userName;
        private String userLevel;
        // 省略getter、setter、构造方法
    }

    // 自定义实体类:富化后的订单
    public static class EnrichedOrder {
        private String orderId;
        private String userId;
        private double amount;
        private String userName;
        private String userLevel;
        // 省略getter、setter、构造方法、toString
    }

    // 自定义Kafka反序列化Schema,比如JSON格式实现
    // class OrderDeserializationSchema implements DeserializationSchema<Order> {...}
    // class UserInfoDeserializationSchema implements DeserializationSchema<UserInfo> {...}
}

如果你更习惯用SQL来处理流数据,Flink SQL的Lookup Join完美适配这个场景。你可以把静态Kafka流定义为一个Lookup Table(配置缓存规则让数据驻留内存),然后通过SQL语句直接关联实时流和Lookup表,代码更简洁。

代码示例

import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.table.api.EnvironmentSettings;
import org.apache.flink.table.api.Table;
import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;

public class SqlLookupEnrichment {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        EnvironmentSettings settings = EnvironmentSettings.newInstance().inStreamingMode().build();
        StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env, settings);

        // 1. 创建实时订单流表
        tableEnv.executeSql("""
                CREATE TABLE real_time_orders (
                    order_id STRING,
                    user_id STRING,
                    amount DOUBLE,
                    ts TIMESTAMP(3) METADATA FROM 'timestamp'
                ) WITH (
                    'connector' = 'kafka',
                    'topic' = 'real-time-orders',
                    'properties.bootstrap.servers' = 'localhost:9092',
                    'properties.group.id' = 'flink-sql-lookup-group',
                    'format' = 'json',
                    'scan.startup.mode' = 'latest-offset'
                )
                """);

        // 2. 创建静态用户Lookup表(配置缓存让数据驻留内存)
        tableEnv.executeSql("""
                CREATE TABLE static_user_info (
                    user_id STRING PRIMARY KEY NOT ENFORCED,
                    user_name STRING,
                    user_level STRING
                ) WITH (
                    'connector' = 'kafka',
                    'topic' = 'static-user-info',
                    'properties.bootstrap.servers' = 'localhost:9092',
                    'format' = 'json',
                    'scan.startup.mode' = 'earliest-offset',
                    -- 配置Lookup缓存:最大缓存10000条,永久有效(如果静态数据不更新,TTL设为0)
                    'lookup.cache.max-rows' = '10000',
                    'lookup.cache.ttl' = '0'
                )
                """);

        // 3. 执行Lookup Join,富化实时订单
        Table enrichedOrders = tableEnv.sqlQuery("""
                SELECT
                    o.order_id,
                    o.user_id,
                    o.amount,
                    u.user_name,
                    u.user_level
                FROM real_time_orders o
                LEFT JOIN static_user_info FOR SYSTEM_TIME AS OF o.ts u
                ON o.user_id = u.user_id
                """);

        // 4. 将结果转成DataStream输出(或者直接写到下游存储)
        tableEnv.toDataStream(enrichedOrders).print();

        env.execute("SQL Lookup Enrichment Job");
    }
}

方案选择建议

  • Broadcast State:适合需要自定义关联逻辑、或者静态数据可能后续有更新的场景,灵活性更高。
  • Flink SQL Lookup Join:适合SQL熟悉的开发者,代码更简洁,无需手动管理状态。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.28 22:22:45