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

搭建Kafka-Flink Connector项目并提供含pom.xml的示例程序

Hey there! 作为刚摸Kafka-Flink Connector的新手,我太懂你想要一个能直接跑通的简单示例的需求了。下面我给你整了一套完整的代码,包括Kafka生产者(往主题写数据)、Flink消费者(监听主题并打印数据),还有对应的pom.xml配置,都是最基础的入门版本,帮你快速走通整个流程。

先做前置准备

  • 确保你本地/服务器上已经启动了Kafka和ZooKeeper(如果是Kafka 2.8及以上版本,也可以用KRaft模式不用ZooKeeper,但新手用传统ZooKeeper模式更稳妥)
  • 提前创建测试主题test-flink-topic,用Kafka自带的命令行工具执行:
    kafka-topics.sh --create --topic test-flink-topic --bootstrap-server localhost:9092 --partitions 1 --replication-factor 1
    

1. Kafka生产者代码(Java)

这个类会循环往test-flink-topic发送带时间戳的消息,方便你测试:

import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.serialization.StringSerializer;

import java.util.Properties;
import java.util.concurrent.TimeUnit;

public class SimpleKafkaProducer {
    public static void main(String[] args) throws InterruptedException {
        // 配置生产者参数
        Properties props = new Properties();
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());

        // 创建生产者实例
        try (KafkaProducer<String, String> producer = new KafkaProducer<>(props)) {
            // 循环发送10条测试消息
            for (int i = 0; i < 10; i++) {
                String message = "Test message " + i + " at " + System.currentTimeMillis();
                ProducerRecord<String, String> record = new ProducerRecord<>("test-flink-topic", message);
                producer.send(record);
                System.out.println("Sent: " + message);
                TimeUnit.SECONDS.sleep(1);
            }
        }
    }
}

2. Flink消费者代码(Java)

这个Flink程序会监听test-flink-topic,把收到的消息直接打印到控制台:

import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.connector.kafka.source.KafkaSource;
import org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;

public class SimpleFlinkKafkaConsumer {
    public static void main(String[] args) throws Exception {
        // 获取Flink执行环境
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        // 为了方便测试,设置并行度为1
        env.setParallelism(1);

        // 创建Kafka Source
        KafkaSource<String> kafkaSource = KafkaSource.<String>builder()
                .setBootstrapServers("localhost:9092")
                .setTopics("test-flink-topic")
                .setGroupId("flink-test-group")
                // 从最早的偏移量开始消费,适合第一次测试
                .setStartingOffsets(OffsetsInitializer.earliest())
                .setValueOnlyDeserializer(new org.apache.flink.api.common.serialization.SimpleStringSchema())
                .build();

        // 读取Kafka数据,转为DataStream
        DataStream<String> kafkaStream = env.fromSource(kafkaSource, WatermarkStrategy.noWatermarks(), "Kafka Source");

        // 打印数据到控制台
        kafkaStream.print("Received from Kafka: ");

        // 执行Flink程序
        env.execute("Flink Kafka Consumer Demo");
    }
}

3. 对应的pom.xml配置

这里用的是Flink 1.17.1和Kafka 3.0.0的版本,版本兼容性比较稳定,你可以根据自己的环境调整:

<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
         xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
    <modelVersion>4.0.0</modelVersion>

    <groupId>com.example</groupId>
    <artifactId>flink-kafka-demo</artifactId>
    <version>1.0-SNAPSHOT</version>

    <properties>
        <flink.version>1.17.1</flink.version>
        <kafka.version>3.0.0</kafka.version>
        <java.version>1.8</java.version>
        <maven.compiler.source>${java.version}</maven.compiler.source>
        <maven.compiler.target>${java.version}</maven.compiler.target>
    </properties>

    <dependencies>
        <!-- Flink核心依赖 -->
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-streaming-java</artifactId>
            <version>${flink.version}</version>
            <scope>provided</scope>
        </dependency>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-clients</artifactId>
            <version>${flink.version}</version>
            <scope>provided</scope>
        </dependency>

        <!-- Flink Kafka Connector -->
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-connector-kafka</artifactId>
            <version>${flink.version}</version>
        </dependency>

        <!-- Kafka客户端依赖 -->
        <dependency>
            <groupId>org.apache.kafka</groupId>
            <artifactId>kafka-clients</artifactId>
            <version>${kafka.version}</version>
        </dependency>
    </dependencies>

    <build>
        <plugins>
            <!-- 打包插件,方便提交到Flink集群运行 -->
            <plugin>
                <groupId>org.apache.maven.plugins</groupId>
                <artifactId>maven-shade-plugin</artifactId>
                <version>3.2.4</version>
                <executions>
                    <execution>
                        <phase>package</phase>
                        <goals>
                            <goal>shade</goal>
                        </goals>
                        <configuration>
                            <artifactSet>
                                <excludes>
                                    <exclude>org.apache.flink:force-shading</exclude>
                                    <exclude>com.google.code.findbugs:jsr305</exclude>
                                    <exclude>org.slf4j:*</exclude>
                                    <exclude>log4j:*</exclude>
                                </excludes>
                            </artifactSet>
                            <filters>
                                <filter>
                                    <artifact>*:*</artifact>
                                    <excludes>
                                        <exclude>META-INF/*.SF</exclude>
                                        <exclude>META-INF/*.DSA</exclude>
                                        <exclude>META-INF/*.RSA</exclude>
                                    </excludes>
                                </filter>
                            </filters>
                            <transformers>
                                <transformer implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer">
                                    <mainClass>com.example.SimpleFlinkKafkaConsumer</mainClass>
                                </transformer>
                            </transformers>
                        </configuration>
                    </execution>
                </executions>
            </plugin>
        </plugins>
    </build>
</project>

运行步骤提示

  1. 先启动ZooKeeper和Kafka服务
  2. 运行Flink消费者程序,让它先监听主题
  3. 运行Kafka生产者程序,开始发送消息
  4. 你就能在Flink程序的控制台看到收到的消息啦!

如果是在Flink集群上运行,记得用mvn package打包成jar包,然后提交到集群执行~

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.06 10:17:39