请求提供基于Kafka与Flink创建连续查询的入门示例
Kafka + Flink 连续查询入门示例
环境准备
- JDK 8 或 11(Flink 推荐适配版本)
- Flink 1.17.x 及以上版本
- Kafka 2.8.x 及以上版本
- Maven/Gradle 构建工具
核心依赖(Maven)
<dependencies> <!-- Flink 核心依赖 --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-java</artifactId> <version>1.17.1</version> <scope>provided</scope> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-java</artifactId> <version>1.17.1</version> <scope>provided</scope> </dependency> <!-- Flink-Kafka 连接器 --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-kafka</artifactId> <version>1.17.1</version> </dependency> </dependencies>
完整代码示例
这个示例实现从Kafka Topic读取用户行为消息,连续统计每个用户的实时点击次数,属于典型的无界流连续查询场景。
import org.apache.flink.api.common.functions.MapFunction; import org.apache.flink.api.common.serialization.SimpleStringSchema; 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.connectors.kafka.FlinkKafkaConsumer; import java.util.Properties; public class KafkaFlinkContinuousQueryDemo { public static void main(String[] args) throws Exception { // 1. 初始化Flink流执行环境 StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(1); // 入门阶段设为1,方便观察结果 // 2. 配置Kafka消费者参数 Properties kafkaProps = new Properties(); kafkaProps.setProperty("bootstrap.servers", "localhost:9092"); // 替换为你的Kafka地址 kafkaProps.setProperty("group.id", "flink-kafka-demo-group"); kafkaProps.setProperty("auto.offset.reset", "latest"); // 从最新消息开始消费 // 3. 从指定Kafka Topic读取数据 String inputTopic = "user-behavior-topic"; DataStream<String> kafkaSource = env.addSource( new FlinkKafkaConsumer<>(inputTopic, new SimpleStringSchema(), kafkaProps) ); // 4. 数据转换:将消息解析为(user_id, 1)二元组 DataStream<Tuple2<String, Integer>> userClicks = kafkaSource.map(new MapFunction<String, Tuple2<String, Integer>>() { @Override public Tuple2<String, Integer> map(String value) throws Exception { // 假设消息格式为:user_id:u101,behavior:click String userId = value.split(",")[0].split(":")[1]; return new Tuple2<>(userId, 1); } }); // 5. 连续查询:按用户分组,实时累计点击次数 DataStream<Tuple2<String, Integer>> resultStream = userClicks .keyBy(tuple -> tuple.f0) // 按user_id分组 .sum(1); // 对点击数做累加,实时更新结果 // 6. 将连续查询结果输出到控制台 resultStream.print("实时点击统计"); // 启动Flink任务 env.execute("Kafka-Flink 连续查询示例"); } }
运行步骤
- 启动Kafka服务:启动ZooKeeper和Kafka,创建指定的
user-behavior-topicTopic - 发送测试消息:用Kafka命令行生产者发送符合格式的消息,比如:
bin/kafka-console-producer.sh --broker-list localhost:9092 --topic user-behavior-topic # 输入测试消息 user_id:u101,behavior:click user_id:u101,behavior:click user_id:u102,behavior:click - 运行Flink程序:在IDE中直接运行主类,或者打包成Jar提交到Flink集群
- 观察结果:控制台会实时输出每个用户的累计点击次数,每新增一条消息,对应用户的统计结果就会更新
关键说明
- 这里的连续查询是无窗口的实时聚合,结果会随着新消息持续更新;如果需要周期性输出结果,可以改用滚动/滑动窗口
- 生产环境中可以将结果输出到Kafka、数据库等持久化存储,替换示例中的控制台输出
- 如果需要处理乱序数据,可以添加水印(Watermark)配置,保证统计的准确性
内容的提问来源于stack exchange,提问作者user20740519
相关产品推荐
相关产品推荐

