【问题标题】:Deserialise a POJO in Kafka Streams反序列化 Kafka Streams 中的 POJO
【发布时间】:2019-01-15 12:02:44
【问题描述】:

我的 Kafka 主题有这种格式的消息

user1,subject1,80|user1,subject2,90 

user2,subject1,70|user2,subject2,100 

and so on. 

我已经创建了如下的用户 POJO。

class User implements Serializable{
/**
 * 
 */
private static final long serialVersionUID = -253687203767610477L;
private String userId;
private String subject;
private String marks;

public User(String userId, String subject, String marks) {
    super();
    this.userId = userId;
    this.subject = subject;
    this.marks = marks;
}

public String getUserId() {
    return userId;
}

public void setUserId(String userId) {
    this.userId = userId;
}
public String getSubject() {
    return subject;
}
public void setSubject(String subject) {
    this.subject = subject;
}
public String getMarks() {
    return marks;
}
public void setMarks(String marks) {
    this.marks = marks;
}
}

我还创建了默认键值序列化

streamProperties.put(
            StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
streamProperties.put(
            StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());

我正在尝试按用户 ID 查找计数,如下所示。我还需要用户对象来执行其他一些功能。

KTable<String, Long> wordCount = streamInput

    .flatMap(new KeyValueMapper<String, String, Iterable<KeyValue<String,User>>>() {

        @Override
        public Iterable<KeyValue<String, User>> apply(String key, String value) {
            String[] userObjects = value.split("|");
            List<KeyValue<String, User>> userList = new LinkedList<>();
            for(String userObject: userObjects) {
                String[] userData = userObject.split(",");
                userList.add(KeyValue.pair(userData[0],
                        new User(userData[0],userData[1],userData[2])));


            }
            return userList;
        }
    })

.groupByKey()
.count();

我收到以下错误

Caused by: org.apache.kafka.streams.errors.StreamsException: A serializer (key: org.apache.kafka.common.serialization.StringSerializer / value: org.apache.kafka.common.serialization.StringSerializer) is not compatible to the actual key or value type (key type: java.lang.String / value type: com.example.testing.dao.User). Change the default Serdes in StreamConfig or provide correct Serdes via method parameters.

我想我需要为用户类提供正确的Serde

【问题讨论】:

  • 您需要定义自己的反序列化器类。请展示你在这方面的尝试。还要找出生产者使用的序列化程序在哪里定义
  • class UserDeserializer implements Deserializer 将是一个好的开始。
  • 你确定你的记录真的是一个字符串中的两个用户对象吗?你能用最多 5 条消息显示控制台消费者输出吗?
  • @cricket_007 kafka 中的每条消息都不会包含许多用户对象的信息,这些用户对象由管道分隔符分隔。每个用户信息以逗号分隔。请检查构造函数和每个用户消息。这就是我使用 flatMap 的原因
  • 我认为问题在于您生成消息并从 Kafka 队列中读取消息的方式。两者都需要以相同的方式进行序列化。尝试在 String Serializable 转换所有消息,然后运行您的代码。它应该工作。如果可行,请尝试将其更改为正确的 JSON 序列化程序,然后读取它。

标签: java apache-kafka pojo apache-kafka-streams


【解决方案1】:

问题在于 Value Serdes。

函数groupBy有两个版本:

  • KStream::KGroupedStream&lt;K, V&gt; groupByKey();
  • KStream::KGroupedStream&lt;K, V&gt; groupByKey(final Grouped&lt;K, V&gt; grouped);

第一个版本在后台调用第二个 Grouped 和默认 Serdes(在您的情况下,它用于键和值 StringSerde

您的flatMap 将消息映射到KeyValue&lt;String, User&gt; 类型,因此值的类型为User

在您的情况下,解决方案将改为使用 groupByKey() 调用 groupByKey(Grouped.with(keySerde, valSerde));,并使用适当的 Serdes。

【讨论】:

  • 在这种情况下,什么是合适的 valSerde?​​span>
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2023-03-23
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-10-14
相关资源
最近更新 更多