如何在Flink中连接两个Kafka Streams?利用静态主表流丰富实时流并实现内存永久Lookup Table
适合你的Flink API方案及代码示例
嘿,这个场景我刚好实操过!针对你想用静态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> {...} }
二、Flink SQL API:Lookup Join 方案
如果你更习惯用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
相关产品推荐
相关产品推荐

