6

我开始玩卡夫卡。我已经设置了一个 zookeeper 配置,并且我设法发送和使用 String 消息。现在我试图传递一个对象(在java中),但由于某种原因,在消费者中解析消息时,我遇到了标题问题。我尝试了几个序列化选项(使用解码器/编码器),并且所有的都返回相同的标头问题。

这是我的代码制作人:

        Properties props = new Properties();
        props.put("zk.connect", "localhost:2181");
        props.put("serializer.class", "com.inneractive.reporter.kafka.EventsDataSerializer");
        ProducerConfig config = new ProducerConfig(props);
        Producer<Long, EventDetails> producer = new Producer<Long, EventDetails>(config);
        ProducerData<Long, EventDetails> data = new ProducerData<Long, EventDetails>("test3", 1, Arrays.asList(new EventDetails());
        try {
           producer.send(data);
        } finally {
           producer.close();
        }

和消费者:

        Properties props = new Properties();
        props.put("zk.connect", "localhost:2181");
        props.put("zk.connectiontimeout.ms", "1000000");
        props.put("groupid", "test_group");

        // Create the connection to the cluster
        ConsumerConfig consumerConfig = new ConsumerConfig(props);
        ConsumerConnector consumerConnector = Consumer.createJavaConsumerConnector(consumerConfig);

        // create 4 partitions of the stream for topic “test”, to allow 4 threads to consume
        Map<String, List<KafkaMessageStream<EventDetails>>> topicMessageStreams =
                consumerConnector.createMessageStreams(ImmutableMap.of("test3", 4), new EventsDataSerializer());
        List<KafkaMessageStream<EventDetails>> streams = topicMessageStreams.get("test3");

        // create list of 4 threads to consume from each of the partitions
        ExecutorService executor = Executors.newFixedThreadPool(4);

        // consume the messages in the threads
        for (final KafkaMessageStream<EventDetails> stream: streams) {
            executor.submit(new Runnable() {
                public void run() {
                    for(EventDetails event: stream) {
                        System.err.println("********** Got message" + event.toString());        
                    }
                }
            });
        }

和我的序列化器:

public  class EventsDataSerializer implements Encoder<EventDetails>, Decoder<EventDetails> {
    public Message toMessage(EventDetails eventDetails) {
        try {
            ObjectMapper mapper = new ObjectMapper(new SmileFactory());
            byte[] serialized = mapper.writeValueAsBytes(eventDetails);
            return new Message(serialized);
} catch (IOException e) {
            e.printStackTrace();
            return null;   // TODO
        }
}
    public EventDetails toEvent(Message message) {
        EventDetails event = new EventDetails();

        ObjectMapper mapper = new ObjectMapper(new SmileFactory());
        try {
            //TODO handle error
            return mapper.readValue(message.payload().array(), EventDetails.class);
        } catch (IOException e) {
            e.printStackTrace();
            return null;
        }

    }
}

这是我得到的错误:

org.codehaus.jackson.JsonParseException: Input does not start with Smile format header (first byte = 0x0) and parser has REQUIRE_HEADER enabled: can not parse
 at [Source: N/A; line: -1, column: -1]

当我与MessagePacka一起工作时,ObjectOutputStream我遇到了一个类似的标题问题。我还尝试将有效负载 CRC32 添加到消息中,但这也无济于事。

我在这里做错了什么?

4

2 回答 2

3

嗯,我没有遇到与您遇到的相同的标头问题,但是当我没有VerifiableProperties在我的编码器/解码器中提供构造函数时,我的项目没有正确编译。不过,缺少的构造函数会破坏杰克逊的反序列化似乎很奇怪。

也许尝试拆分您的编码器和解码器,并VerifiableProperties在两者中包含构造函数;您不需要实现Decoder[T]序列化。我能够按照本文中的格式成功实现 jsonObjectMapper序列化。

祝你好运!

于 2014-06-04T17:39:02.597 回答
1

Bytebuffers .array() 方法不是很可靠。这取决于特定的实现。你可能想试试

ByteBuffer bb = message.payload()

byte[] b = new byte[bb.remaining()]
bb.get(b, 0, b.length);
return mapper.readValue(b, EventDetails.class) 
于 2012-08-10T08:39:19.150 回答