INTEXT-84 Kafka: Enhance Avro serialization support

This commit is contained in:
Soby Chacko
2013-07-16 22:40:28 -04:00
committed by Gunnar Hillert
parent d99de8c740
commit b28bf4b06c
38 changed files with 673 additions and 386 deletions

View File

@@ -19,6 +19,7 @@ import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.springframework.context.support.ClassPathXmlApplicationContext;
import org.springframework.integration.MessageChannel;
import org.springframework.integration.samples.kafka.user.User;
import org.springframework.integration.support.MessageBuilder;
public class OutboundRunner {
@@ -34,8 +35,11 @@ public class OutboundRunner {
//sending 100,000 messages to Kafka server for topic test1
for (int i = 0; i < 500; i++) {
final User user = new User();
user.setFirstName("fname" + i);
user.setLastName("lname" + i);
channel.send(
MessageBuilder.withPayload("hello Fom ob adapter test1 - " + i)
MessageBuilder.withPayload(user)
.setHeader("messageKey", String.valueOf(i))
.setHeader("topic", "test1").build());

View File

@@ -27,22 +27,26 @@
<int:poller fixed-delay="1" time-unit="MILLISECONDS"/>
</int-kafka:inbound-channel-adapter>
<bean id="kafkaDecoder" class="org.springframework.integration.kafka.serializer.avro.AvroBackedKafkaDecoder">
<bean id="kafkaReflectionDecoder" class="org.springframework.integration.kafka.serializer.avro.AvroReflectDatumBackedKafkaDecoder">
<constructor-arg type="java.lang.Class" value="java.lang.String"/>
</bean>
<bean id="kafkaSpecificDecoder" class="org.springframework.integration.kafka.serializer.avro.AvroSpecificDatumBackedKafkaDecoder">
<constructor-arg value="org.springframework.integration.samples.kafka.user.User" />
</bean>
<int-kafka:consumer-context id="consumerContext"
consumer-timeout="1000"
zookeeper-connect="zookeeperConnect">
<int-kafka:consumer-configurations>
<int-kafka:consumer-configuration group-id="default"
value-decoder="kafkaDecoder"
key-decoder="kafkaDecoder"
value-decoder="kafkaSpecificDecoder"
key-decoder="kafkaReflectionDecoder"
max-messages="5000">
<int-kafka:topic id="test1" streams="4"/>
</int-kafka:consumer-configuration>
<int-kafka:consumer-configuration group-id="default1"
max-messages="5">
max-messages="50">
<int-kafka:topic id="test2" streams="4"/>
</int-kafka:consumer-configuration>
</int-kafka:consumer-configurations>

View File

@@ -21,20 +21,24 @@
<task:executor id="taskExecutor" pool-size="5" keep-alive="120" queue-capacity="500"/>
<bean id="kafkaEncoder" class="org.springframework.integration.kafka.serializer.avro.AvroBackedKafkaEncoder">
<bean id="kafkaReflectionEncoder" class="org.springframework.integration.kafka.serializer.avro.AvroReflectDatumBackedKafkaEncoder">
<constructor-arg value="java.lang.String" />
</bean>
<bean id="kafkaSpecificEncoder" class="org.springframework.integration.kafka.serializer.avro.AvroSpecificDatumBackedKafkaEncoder">
<constructor-arg value="org.springframework.integration.samples.kafka.user.User" />
</bean>
<bean id="customPartitioner" class="org.springframework.integration.samples.kafka.outbound.CustomPartitioner"/>
<int-kafka:producer-context id="kafkaProducerContext">
<int-kafka:producer-configurations>
<int-kafka:producer-configuration broker-list="localhost:9092"
key-class-type="java.lang.String"
value-class-type="java.lang.String"
value-class-type="org.springframework.integration.samples.kafka.user.User"
topic="test1"
value-encoder="kafkaEncoder"
key-encoder="kafkaEncoder"
value-encoder="kafkaSpecificEncoder"
key-encoder="kafkaReflectionEncoder"
compression-codec="default"
partitioner="customPartitioner"/>
<int-kafka:producer-configuration broker-list="localhost:9092"