惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

J
Java Code Geeks
腾讯CDC
博客园 - 聂微东
爱范儿
爱范儿
罗磊的独立博客
P
Proofpoint News Feed
博客园 - Franky
博客园 - 三生石上(FineUI控件)
OSCHINA 社区最新新闻
OSCHINA 社区最新新闻
酷 壳 – CoolShell
酷 壳 – CoolShell
Jina AI
Jina AI
Blog — PlanetScale
Blog — PlanetScale
让小产品的独立变现更简单 - ezindie.com
让小产品的独立变现更简单 - ezindie.com
博客园 - 司徒正美
美团技术团队
MongoDB | Blog
MongoDB | Blog
WordPress大学
WordPress大学
A
About on SuperTechFans
I
InfoQ
博客园_首页
钛媒体:引领未来商业与生活新知
钛媒体:引领未来商业与生活新知
H
Help Net Security
Microsoft Azure Blog
Microsoft Azure Blog
G
Google Developers Blog

博客园 - 刚泡

数据库短暂波动导致xxl-job-admin不调度任务的问题 CompletableFuture用法详解 使用责任链模式简化if-else代码示例 使用Function Interface简化if-else代码示例 自定义AngularJS Modal弹框大小的方法 SpringCloud Alibaba 改造Sentinel Dashboard将熔断规则持久化到Nacos SpringCloud Alibaba 改造Sentinel Dashboard将流控规则持久化到Nacos [Access][Microsoft][ODBC 驱动程序管理器] 无效的字符串或缓冲区长度 Invalid string or buffer length pom.xml错误:org.codehaus.plexus.archiver.jar.Manifest.write(java.io.PrintWriter)的解决方法 Spring Cloud @HystrixCommand和@CacheResult注解使用,参数配置 Spring Boot的属性加载顺序 Spring Boot 使用properties如何多环境配置 Spring Cloud微服务实战阅读笔记(一) 基础知识 Java使用POI为Excel打水印,调整列宽并设置Excel只读(用户不可编辑) Sitemesh 3使用及配置 Java上传文件夹(Jersey) 路径名导致的异常:javax.imageio.IIOException: Can't read input file! 解决安卓微信浏览器中上传不能调用相机的问题 Mac下安装SVN插件javaHL not available的解决方法
Java操作Kafka执行不成功的解决方法,Kafka Broker Advertise...
刚泡 · 2018-04-21 · via 博客园 - 刚泡

创建Spring Boot项目继承Kafka,向Kafka发送消息始终不成功。具体项目配置如下:

<?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.0http://maven.apache.org/xsd/maven-4.0.0.xsd">

    <modelVersion>4.0.0</modelVersion>

    <groupId>com.lodestone</groupId>

    <artifactId>lodestone-kafka</artifactId>

    <version>0.0.1-SNAPSHOT</version>

    <packaging>jar</packaging>

    <name>lodestone-kafka</name>

    <description>Lodestone kafka</description>

    <parent>

         <groupId>org.springframework.boot</groupId>

        <artifactId>spring-boot-starter-parent</artifactId>

         <version>1.5.12.RELEASE</version>

         <relativePath />

    </parent>

    <properties>

        <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>

        <project.reporting.outputEncoding>UTF-8</project.reporting.outputEncoding>

         <java.version>1.8</java.version>

    </properties>

    <dependencies>

         <dependency>

            <groupId>org.springframework.kafka</groupId>

           <artifactId>spring-kafka</artifactId>

         </dependency>

         <dependency>

             <groupId>org.springframework.boot</groupId>

             <artifactId>spring-boot-starter</artifactId>

         </dependency>

         <dependency>

             <groupId>org.springframework.boot</groupId>

             <artifactId>spring-boot-starter-test</artifactId>

         </dependency>

         <!--  https://mvnrepository.com/artifact/com.alibaba/fastjson -->

         <dependency>

             <groupId>com.alibaba</groupId>

             <artifactId>fastjson</artifactId>

             <version>1.2.47</version>

         </dependency>

         <!--  https://mvnrepository.com/artifact/org.slf4j/slf4j-api -->

         <dependency>

             <groupId>org.slf4j</groupId>

             <artifactId>slf4j-api</artifactId>

         </dependency>

         <!--  https://mvnrepository.com/artifact/org.slf4j/log4j-over-slf4j -->

         <dependency>

             <groupId>org.slf4j</groupId>

              <artifactId>log4j-over-slf4j</artifactId>

         </dependency>

    </dependencies>

    <build>

         <plugins>

             <plugin>

                 <groupId>org.springframework.boot</groupId>

                 <artifactId>spring-boot-maven-plugin</artifactId>

             </plugin>

         </plugins>

    </build>

</project>

application.yml配置:

spring:

  kafka:

    bootstrap-servers:

      - 192.168.52.131:9092

    consumer:

      auto-offset-reset: earliest

      group-id: console-consumer-53989

      key-deserializer:org.apache.kafka.common.serialization.StringDeserializer

      value-deserializer:org.apache.kafka.common.serialization.StringDeserializer

    producer:

      key-serializer:org.apache.kafka.common.serialization.StringSerializer

      value-serializer:org.apache.kafka.common.serialization.StringSerializer

logging:

  level:

    root: DEBUG

    org:

      springframework: DEBUG

      mybatis: DEBUG

生产者代码:

package com.lodestone.kafka.producer;

import java.util.Date;

import java.util.UUID;

import org.springframework.beans.factory.annotation.Autowired;

import org.springframework.kafka.core.KafkaTemplate;

import com.alibaba.fastjson.JSON;

import com.lodestone.kafka.message.LodestoneMessage;

@Component

public class Sender {

    @Autowired

    private KafkaTemplate kafkaTemplate;

    public void sendMessage() {

        LodestoneMessage message = new LodestoneMessage();

        message.setId(UUID.randomUUID().toString().replaceAll("-", ""));

        message.setMsg(UUID.randomUUID().toString());

        message.setSendTime(new Date());

        kafkaTemplate.send("test", JSON.toJSONString(message));

    }

}

消息代码定义:

package com.lodestone.kafka.message;

import java.io.Serializable;

import java.util.Date;

public class LodestoneMessage implements Serializable {

    private static final long serialVersionUID = -6847574917429814430L;

    private String id;

    private String msg;

    private Date sendTime;

    public String getId() {

        return id;

    }

    public void setId(String id) {

        this.id = id;

    }

    public String getMsg() {

        return msg;

    }

    public void setMsg(String msg) {

        this.msg = msg;

    }

    public Date getSendTime() {

        return sendTime;

    }

    public void setSendTime(Date sendTime) {

        this.sendTime = sendTime;

    }

}

Spring Boot应用启动类:

package com.lodestone.kafka;

import org.springframework.boot.SpringApplication;

import org.springframework.boot.autoconfigure.SpringBootApplication;

import org.springframework.context.ApplicationContext;

import com.lodestone.kafka.producer.Sender;

@SpringBootApplication

public class LodestoneKafkaApplication {    

    public static void main(String[] args) throws InterruptedException {

        ApplicationContext app = SpringApplication.run(LodestoneKafkaApplication.class, args);

        //测试

        for(int i=0; i<5; i++) {

            Sender sender = app.getBean(Sender.class);

            sender.sendMessage();

            Thread.sleep(500);

        }

    }

}

将应用的日志调整位Debug级别,启动应用时看到如下报错:

    at sun.nio.ch.SocketChannelImpl.checkConnect(Native Method) ~[na:1.8.0_111]

    at sun.nio.ch.SocketChannelImpl.finishConnect(Unknown Source) ~[na:1.8.0_111]

    at org.apache.kafka.clients.producer.internals.Sender.run(Sender.java:236) [kafka-clients-0.10.1.1.jar:na]

    at org.apache.kafka.clients.producer.internals.Sender.run(Sender.java:148) [kafka-clients-0.10.1.1.jar:na]

    at java.lang.Thread.run(Unknown Source) [na:1.8.0_111]

可以看到报错第一句显示:Connection with localhost/127.0.0.1 disconnected

但是我明明在application.yml中配置了我的Kafka Server的地址是:192.168.52.131:9092,而在实际连接kafka服务器时却使用的localhost/127.0.0.1这个地址,所以导致无法连接kafka Server。

经过百度,得知,在设置Kafka的时候,需要设置advertised.listeners这个属性。

该属性在config/server.properties中的描述如下:

# Hostname and port the broker will advertise to producers and consumers. If not set,

# it uses the value for "listeners" if configured.  Otherwise, it will use the value

# advertised.listeners=PLAINTEXT://:your.host.name:9092

"PLAINTEXT"表示协议,可选的值有PLAINTEXT和SSL,hostname可以指定IP地址,也可以用"0.0.0.0"表示对所有的网络接口有效,如果hostname为空表示只对默认的网络接口有效

也就是说如果你没有配置advertised.listeners,就使用listeners的配置通告给消息的生产者和消费者,这个过程是在生产者和消费者获取源数据(metadata)。

因此重新设置advertised.listeners为如下:

advertised.listeners=PLAINTEXT://192.168.52.131:9092

需要注意的是,如果Kafka有多个节点,那么需要每个节点都按照这个节点的实际hostname和port情况进行设置。

每个节点都设置之后,再重新启动Spring Boot应用,则能够正常连接Kafka Server,并能够正常发送消息了。

部分援引和参考: