搭建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>
运行步骤提示
- 先启动ZooKeeper和Kafka服务
- 运行Flink消费者程序,让它先监听主题
- 运行Kafka生产者程序,开始发送消息
- 你就能在Flink程序的控制台看到收到的消息啦!
如果是在Flink集群上运行,记得用mvn package打包成jar包,然后提交到集群执行~
内容的提问来源于stack exchange,提问作者Noobie
相关产品推荐
相关产品推荐

