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

请求提供基于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 连续查询示例");
    }
}

运行步骤

  1. 启动Kafka服务:启动ZooKeeper和Kafka,创建指定的user-behavior-topic Topic
  2. 发送测试消息:用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
    
  3. 运行Flink程序:在IDE中直接运行主类,或者打包成Jar提交到Flink集群
  4. 观察结果:控制台会实时输出每个用户的累计点击次数,每新增一条消息,对应用户的统计结果就会更新

关键说明

  • 这里的连续查询是无窗口的实时聚合,结果会随着新消息持续更新;如果需要周期性输出结果,可以改用滚动/滑动窗口
  • 生产环境中可以将结果输出到Kafka、数据库等持久化存储,替换示例中的控制台输出
  • 如果需要处理乱序数据,可以添加水印(Watermark)配置,保证统计的准确性

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 14:35:10