【发布时间】:2015-08-11 16:07:27
【问题描述】:
我正在尝试通过我的演员向 Kafka 发送消息,但它不起作用。
以下代码有效
new KeyedMessage[String, Array[Byte]]("my-topic", msg.message)
这个不...为什么?
new KeyedMessage[String, Array[Byte]]("my-topic", msg.id, msg.message)
甚至
new KeyedMessage[String, Array[Byte]]("my-topic", msg.id, null, msg.message)
将partition设置为null,强制只填充id,但还是一样。
有什么想法吗?
编辑
好吧,只是意识到当以字节数组发送消息时它不起作用。
我已经更改了我的代码以允许以字符串形式发送消息并且它可以工作。
以下内容不起作用:
props.put("serializer.class", "kafka.serializer.DefaultEncoder")
override def receive: Receive = {
case msg: Message =>
val keyedMessage = new KeyedMessage[String, Array[Byte]]("my-topic", msg.id, msg.message)
producer.send(keyedMessage)
case _ => log.error("Got a msg that I don't understand")
}
将 Array[Byte] 修改为 String,同时将 DefaultEncoder 修改为 StringEncoder,就可以了。
有什么想法吗?
【问题讨论】:
标签: scala akka apache-kafka