【发布时间】:2019-05-04 15:11:25
【问题描述】:
我在storm中实现了一个Logger Bolt,元组的输入来自Kafka Topic。我正在使用 Kafka Connect 来监听 mySQL 数据库的变化。
public class LoggerBolt extends BaseBasicBolt {
private static final long serialVersionUID = 1L;
private static final Logger LOG = Logger.getLogger(LoggerBolt.class);
public void execute(Tuple input, BasicOutputCollector collector) {
System.out.println(input.getValue(0));
}
public void declareOutputFields(OutputFieldsDeclarer declarer) {
}
}
当在本地集群上运行时,下面会打印出来。
Q�%Buckley, Rose RoseBuckley"BuckleyR@univ.edu"963.555.6855x5018963.777.5233策展人 Q� Stanton, Kathie KathieStanton"StantonK@univ.edu963.555.7095963.777.1015教授 Q银行,香农香农 BanksBanksS@univ.edu963.555.7198963.777.6979教授 Q�/巴恩斯,Cleo CleoBarnes BarnesC@univ.edu"963.555.7463x7335963.777.1583$研究教授
我想将这些细节转换为模型类的 Person 对象? 我们如何将 Tuple 输入解析成一个对象?
我尝试了input.getValues(0) , input.getFields(0) 和其他方法,似乎都没有。
【问题讨论】:
-
input.getValue(0)的类型是什么? (尝试打印 input.getValue(0).getClass()) 是否已经是 Person 或 String 或其他?
-
输入是 TupleImpl input.getValue 是 Avro
-
好的。我认为您需要在写入 Kafka 之前查看如何序列化数据。一旦你知道数据是如何序列化的,我们将更容易告诉你应该如何反序列化它。另请注意您使用的是storm-kafka 还是storm-kafka-client spout。
-
如何在 Storm 中配置 Kafka 反序列化器?还是它总是假设字符串?
标签: java apache-kafka apache-storm apache-kafka-connect