当前位置: 首页 > 知识库问答 >
问题:

为Kafka在《春云流水》中设定基调

方和宜
2023-03-14

我正试图与Kafka建立一个Spring云流项目。除密钥反序列化外,一切都按预期运行。这些文件是pom.xml:

<?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 https://maven.apache.org/xsd/maven-4.0.0.xsd">
    <modelVersion>4.0.0</modelVersion>
    <parent>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-parent</artifactId>
        <version>2.5.8</version>
        <relativePath/> <!-- lookup parent from repository -->
    </parent>
    <groupId>com.example</groupId>
    <artifactId>cloud-stream</artifactId>
    <version>0.0.1-SNAPSHOT</version>
    <name>cloud-stream</name>
    <description>Demo project for Spring Boot</description>
    <properties>
        <java.version>11</java.version>
        <spring-cloud.version>2020.0.5</spring-cloud.version>
    </properties>
    <dependencies>
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-web</artifactId>
        </dependency>
        <dependency>
            <groupId>org.apache.kafka</groupId>
            <artifactId>kafka-streams</artifactId>
        </dependency>
        <dependency>
            <groupId>org.springframework.cloud</groupId>
            <artifactId>spring-cloud-stream</artifactId>
        </dependency>
        <dependency>
            <groupId>org.springframework.cloud</groupId>
            <artifactId>spring-cloud-stream-binder-kafka</artifactId>
        </dependency>
        <dependency>
            <groupId>org.springframework.cloud</groupId>
            <artifactId>spring-cloud-stream-binder-kafka-streams</artifactId>
        </dependency>
        <dependency>
            <groupId>org.springframework.kafka</groupId>
            <artifactId>spring-kafka</artifactId>
        </dependency>

        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-test</artifactId>
            <scope>test</scope>
        </dependency>
        <dependency>
            <groupId>org.springframework.cloud</groupId>
            <artifactId>spring-cloud-stream</artifactId>
            <scope>test</scope>
            <classifier>test-binder</classifier>
            <type>test-jar</type>
        </dependency>
        <dependency>
            <groupId>org.springframework.kafka</groupId>
            <artifactId>spring-kafka-test</artifactId>
            <scope>test</scope>
        </dependency>
    </dependencies>
    <dependencyManagement>
        <dependencies>
            <dependency>
                <groupId>org.springframework.cloud</groupId>
                <artifactId>spring-cloud-dependencies</artifactId>
                <version>${spring-cloud.version}</version>
                <type>pom</type>
                <scope>import</scope>
            </dependency>
        </dependencies>
    </dependencyManagement>

    <build>
        <plugins>
            <plugin>
                <groupId>org.springframework.boot</groupId>
                <artifactId>spring-boot-maven-plugin</artifactId>
            </plugin>
        </plugins>
    </build>

</project>

消费者和主要阶层

package com.example;

import java.util.function.Consumer;

import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.context.annotation.Bean;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageHeaders;

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

@SpringBootApplication
public class App {

    private Logger logger = LoggerFactory.getLogger(App.class);

    public static void main(String[] args) {
        SpringApplication.run(App.class, args);
    }

    @Bean
    public Consumer<Message<ChatMessage>> consumeString(){
        return message -> {
            MessageHeaders headers = message.getHeaders();
            logger.info("sup");
        };
    }

    @Bean
    public Consumer<Message<String>> consumeString2(){
        return message -> {
            MessageHeaders headers = message.getHeaders();
            logger.info("sup");
        };
    }
}

一个控制器,它使用 StreamBridge 按需发送 kafka 消息(注意 streamBridge 正在尝试将消息键标头作为字符串发送)

package com.example;

import org.springframework.cloud.stream.function.StreamBridge;
import org.springframework.kafka.support.KafkaHeaders;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RestController;

@RestController
public class Controller {

    private StreamBridge streamBridge;

    public Controller(StreamBridge streamBridge){
        this.streamBridge = streamBridge;
    }

    @PostMapping("/sendMessage/message")
    public String publishMessageString(@RequestBody ChatMessage payload) {
        streamBridge.send("output-out-0",
            MessageBuilder.withPayload(payload)
                .setHeader(KafkaHeaders.MESSAGE_KEY, "message")
                .build());
        return "Success";
    }

}

模型类

package com.example;

public class ChatMessage {

    private String contents;
    private long time;

    public ChatMessage() {

    }

    public ChatMessage(String contents, long time) {

        this.contents = contents;
        this.time = time;
    }

    public String getContents() {

        return contents;
    }

    public long getTime() {

        return time;
    }

    public void setTime(long time) {

        this.time = time;
    }

    @Override
    public String toString() {

        return "ChatMessage [contents=" + contents + ", time=" + time + "]";
    }

}

最后这是应用程序。yml

spring:
  cloud:
    stream:
      function:
        definition: consumeString;consumeString2
      binders:
        binder1:
          type: kafka
          environment.spring.cloud.stream.kafka.streams.binder:
            brokers:
              - localhost:9092
        binder2:
          type: kafka
          environment.spring.cloud.stream.kafka.streams.binder:
            brokers:
              - localhost:9092
      kafka:
        streams:
          binder:
            configuration:
              default:
                key.serde: org.apache.kafka.common.serialization.Serdes$StringSerde

      bindings:
        consumeString-in-0:
          binder: binder1
          destination: test
          group: input-group-1
        consumeString2-in-0:
          binder: binder2
          destination: test2
          group: input-group-2
        output-out-0:
          binder: binder1
          destination: test

当我尝试通过控制器发送消息时,我收到以下错误:

java.lang.ClassCastException: class java.lang.String cannot be cast to class [B (java.lang.String and [B are in module java.base of loader 'bootstrap')
    at org.apache.kafka.common.serialization.ByteArraySerializer.serialize(ByteArraySerializer.java:19) ~[kafka-clients-2.7.2.jar:na]
    at org.apache.kafka.common.serialization.Serializer.serialize(Serializer.java:62) ~[kafka-clients-2.7.2.jar:na]
    at org.apache.kafka.clients.producer.KafkaProducer.doSend(KafkaProducer.java:918) ~[kafka-clients-2.7.2.jar:na]

问题似乎是 kafka 生产者和消费者没有进行字符串反序列化,而是字节数组反序列化。如何更改默认的反序列化程序?

共有1个答案

阳福
2023-03-14

可以在绑定程序或绑定级别重写默认序列化程序。例如,对于特定绑定:

spring:
  cloud:
    stream:
      kafka:
        bindings:
          output-out-0:
            producer:
              configuration:
                "[key.serializer]": ...

https://docs.spring.io/spring-cloud-stream-binder-kafka/docs/3.2.1/reference/html/spring-cloud-stream-binder-kafka.html#kafka-producer-properties

 类似资料:
  • 我用的是Apache Kafka 2.7.0和Spring Cloud Stream Kafka Streams。 在我的Spring Cloud Stream (Kafka Streams)应用程序中,我已经将我的application.yml配置为当输入主题中的消息出现反序列化错误时使用sendToDlq机制: 我启动了我的应用程序,但我看不到这个主题存在。文档指出,如果 DLQ 主题不存在,

  • Spring Cloud Kafka Streams与Spring Cloud Stream、Spring Cloud Function、Spring AMQP和Spring for Apache Kafka有什么区别?

  • 我正试图按照GitHub的建议设置测试 其中StreamProcessor设置为 -->line从不使用在我看来应该在主题“output”上的消息,因为@StreamProcessor有@Sendto(“output”) 我希望能够测试流处理的消息。

  • 我有一个常见的任务问题,我可以找到任何解决方案或帮助(也许我需要传递一些属性来工作?)我使用本地服务器1.3.0.M2并创建简单的流 在日志中,我得到了这个: 2017-09-28 12:31:00.644 信息 5156 --- [ -C-1] o.. a.k.c.c.internals.AbstractCoordinator : 成功加入第 1 代的组测试 2017-09-28 12:31:0

  • 现在我正在尝试用kafka创建消息服务功能以使用< code > spring-cloud-stream-bind-Kafka ,但效果不太好。 Spring罩1.4.2 当我使用此错误日志启动项目时失败 我在怀疑我的春靴版本。这么低配的版本。< br >我认为< code > spring-cloud-stream-binder-Kafka 在spring boot 2.0版本下无法使用或者其他

  • 我想在我的spring boot项目中使用Kafka Streams实时处理。所以我需要Kafka Streams配置,或者我想使用KStreams或KTable,但我在互联网上找不到示例。 我做了制作人和消费者现在我想流实时。