Flink连接Kafka报Connection refused错误的解决求助
问题:Flink消费者连接Kafka时出现Connection refused错误
Kafka配置(server.properties)
broker.id=0 listeners=PLAINTEXT://localhost:9092 log.dirs=/tmp/kafka-logs zookeeper.connect=localhost:2181 num.network.threads=3 num.io.threads=8 socket.send.buffer.bytes=102400 socket.receive.buffer.bytes=102400 socket.request.max.bytes=104857600 log.retention.hours=168 log.segment.bytes=1073741824 log.retention.check.interval.ms=300000 log.cleaner.enable=true
服务验证情况
已按顺序启动Zookeeper和Kafka,无报错;用nc命令验证端口:
nc -zv localhost 9092
返回:
Connection to localhost (127.0.0.1) 9092 port [tcp/*] succeeded!
Flink程序依赖
<dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-java_2.12</artifactId> <version>1.13.2</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-clients_2.12</artifactId> <version>1.13.2</version> </dependency> <dependency> <groupId>junit</groupId> <artifactId>junit</artifactId> <version>4.13.2</version> <scope>test</scope> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-kafka_2.12</artifactId> <version>1.14.6</version> </dependency> <dependency> <groupId>com.fasterxml.jackson.core</groupId> <artifactId>jackson-databind</artifactId> <version>1.12.3</version> </dependency>
程序代码
import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; import org.apache.flink.api.common.functions.AggregateFunction; 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.api.windowing.time.Time; import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer; import org.apache.flink.api.common.serialization.SimpleStringSchema; import java.util.Properties; import org.apache.flink.api.common.functions.FlatMapFunction; import org.apache.flink.util.Collector; public class WordCount { public static void main(String[] args) throws Exception { // Set up the streaming execution environment final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); // Create a DataStream from a source DataStream<String> text = env.socketTextStream("localhost", 9999); // Process the stream DataStream<Tuple2<String, Integer>> wordCounts = text .flatMap(new Tokenizer()) .keyBy(value -> value.f0) .sum(1); // Print the results to the standard output wordCounts.print(); // Execute the job env.execute("Flink Streaming Word Count"); // Kafka properties Properties properties = new Properties(); properties.setProperty("bootstrap.servers", "localhost:9092"); properties.setProperty("group.id", "test"); // Create a Kafka consumer FlinkKafkaConsumer<String> consumer = new FlinkKafkaConsumer<>( "test", new SimpleStringSchema(), properties ); // Add the consumer as a source DataStream<String> stream = env.addSource(consumer); // Parse the JSON transaction records and extract transaction amounts DataStream<Double> transactionAmounts = stream.map(record -> { ObjectMapper objectMapper = new ObjectMapper(); JsonNode jsonNode = objectMapper.readTree(record); return jsonNode.get("tranamount").asDouble(); }); // Calculate sum, count, and average over a sliding window of 1 minute DataStream<Tuple2<String, Double>> result = transactionAmounts .timeWindowAll(Time.minutes(1)) .aggregate(new AggregateFunction<Double, Tuple2<Double, Integer>, Tuple2<String, Double>>() { @Override public Tuple2<Double, Integer> createAccumulator() { return new Tuple2<>(0.0, 0); } @Override public Tuple2<Double, Integer> add(Double value, Tuple2<Double, Integer> accumulator) { return new Tuple2<>(accumulator.f0 + value, accumulator.f1 + 1); } @Override public Tuple2<String, Double> getResult(Tuple2<Double, Integer> accumulator) { double sum = accumulator.f0; int count = accumulator.f1; double average = sum / count; return new Tuple2<>("Sum: " + sum + ", Count: " + count + ", Avg: " + average, average); } @Override public Tuple2<Double, Integer> merge(Tuple2<Double, Integer> a, Tuple2<Double, Integer> b) { return new Tuple2<>(a.f0 + b.f0, a.f1 + b.f1); } }); // Print the result to the console result.print(); // Execute the Flink job env.execute("Flink Kafka JSON Consumer Example"); } public static final class Tokenizer implements FlatMapFunction<String, Tuple2<String, Integer>> { @Override public void flatMap(String value, Collector<Tuple2<String, Integer>> out) { // Normalize and split the line into words String[] tokens = value.toLowerCase().split("\\W+"); // Emit the pairs for (String token : tokens) { if (token.length() > 0) { out.collect(new Tuple2<>(token, 1)); } } } } }
错误信息
Caused by: java.net.ConnectException: Connection refused (Connection refused)
解决方案
1. 修复程序执行顺序问题
代码中env.execute("Flink Streaming Word Count");会直接启动并阻塞执行第一个Socket流作业,后续的Kafka相关代码永远不会被执行。实际报错的是第一个作业中的env.socketTextStream("localhost", 9999);——因为9999端口没有监听服务,导致连接被拒绝。
解决方法:删除第一个Socket流相关的代码块,或者将两个作业拆分为独立的类执行。修改后的main方法核心逻辑如下:
public static void main(String[] args) throws Exception { final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); // Kafka properties Properties properties = new Properties(); properties.setProperty("bootstrap.servers", "localhost:9092"); properties.setProperty("group.id", "test"); // Create a Kafka consumer FlinkKafkaConsumer<String> consumer = new FlinkKafkaConsumer<>( "test", new SimpleStringSchema(), properties ); // Add the consumer as a source DataStream<String> stream = env.addSource(consumer); // Parse the JSON transaction records and extract transaction amounts DataStream<Double> transactionAmounts = stream.map(record -> { ObjectMapper objectMapper = new ObjectMapper(); JsonNode jsonNode = objectMapper.readTree(record); return jsonNode.get("tranamount").asDouble(); }); // Calculate sum, count, and average over a sliding window of 1 minute DataStream<Tuple2<String, Double>> result = transactionAmounts .timeWindowAll(Time.minutes(1)) .aggregate(new AggregateFunction<Double, Tuple2<Double, Integer>, Tuple2<String, Double>>() { @Override public Tuple2<Double, Integer> createAccumulator() { return new Tuple2<>(0.0, 0); } @Override public Tuple2<Double, Integer> add(Double value, Tuple2<Double, Integer> accumulator) { return new Tuple2<>(accumulator.f0 + value, accumulator.f1 + 1); } @Override public Tuple2<String, Double> getResult(Tuple2<Double, Integer> accumulator) { double sum = accumulator.f0; int count = accumulator.f1; double average = sum / count; return new Tuple2<>("Sum: " + sum + ", Count: " + count + ", Avg: " + average, average); } @Override public Tuple2<Double, Integer> merge(Tuple2<Double, Integer> a, Tuple2<Double, Integer> b) { return new Tuple2<>(a.f0 + b.f0, a.f1 + b.f1); } }); // Print the result to the console result.print(); // Execute the Flink job env.execute("Flink Kafka JSON Consumer Example"); }
2. 依赖版本对齐
Flink核心版本是1.13.2,但Kafka连接器版本是1.14.6,版本不匹配可能引发兼容性问题。建议将连接器版本改为与Flink核心一致:
<dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-kafka_2.12</artifactId> <version>1.13.2</version> </dependency>
3. 确认Kafka Topic存在
确保test Topic已在Kafka中创建,可通过以下命令检查:
kafka-topics.sh --list --bootstrap-server localhost:9092
如果不存在,创建Topic:
kafka-topics.sh --create --topic test --bootstrap-server localhost:9092 --partitions 1 --replication-factor 1
内容的提问来源于stack exchange,提问作者M_Gh
相关产品推荐
相关产品推荐

