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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 05:30:54