GH-84, GH-69 Improve Partitioner deprecation note

Fixes GH-84 (https://github.com/spring-projects/spring-integration-kafka/issues/84)
Fixes GH-69 (https://github.com/spring-projects/spring-integration-kafka/issues/69)

* The deprecation for the `partitioner` option hasn't mentioned the `partition-id` (`partition-id-expression`) option.
* Add built-in conversion for the `String <-> byte[]` to avoid serialization for Strings
* Expose `charset` option to configure `String <-> byte[]` conversion.
* Fix `deprecation` message for the `KafkaConsumerContextParser`

Extract `StringBytesConverter`
This commit is contained in:
Artem Bilan
2015-11-05 21:07:55 -05:00
committed by Artem Bilan
parent 0f756cc1ec
commit ba96e7ce71
10 changed files with 108 additions and 51 deletions

View File

@@ -30,13 +30,6 @@ import org.springframework.beans.factory.support.ManagedMap;
import org.springframework.beans.factory.xml.AbstractSingleBeanDefinitionParser;
import org.springframework.beans.factory.xml.ParserContext;
import org.springframework.integration.config.xml.IntegrationNamespaceUtils;
import org.springframework.integration.kafka.support.ConsumerConfigFactoryBean;
import org.springframework.integration.kafka.support.ConsumerConfiguration;
import org.springframework.integration.kafka.support.ConsumerConnectionProvider;
import org.springframework.integration.kafka.support.ConsumerMetadata;
import org.springframework.integration.kafka.support.KafkaConsumerContext;
import org.springframework.integration.kafka.support.MessageLeftOverTracker;
import org.springframework.integration.kafka.support.TopicFilterConfiguration;
import org.springframework.util.StringUtils;
import org.springframework.util.xml.DomUtils;
@@ -54,7 +47,7 @@ public class KafkaConsumerContextParser extends AbstractSingleBeanDefinitionPars
@Override
protected Class<?> getBeanClass(final Element element) {
return KafkaConsumerContext.class;
return org.springframework.integration.kafka.support.KafkaConsumerContext.class;
}
@Override
@@ -70,9 +63,9 @@ public class KafkaConsumerContextParser extends AbstractSingleBeanDefinitionPars
Map<String, BeanMetadataElement> consumerConfigurationsMap = new ManagedMap<String, BeanMetadataElement>();
for (final Element consumerConfiguration : DomUtils.getChildElementsByTagName(consumerConfigurations, "consumer-configuration")) {
final BeanDefinitionBuilder consumerConfigurationBuilder =
BeanDefinitionBuilder.genericBeanDefinition(ConsumerConfiguration.class);
BeanDefinitionBuilder.genericBeanDefinition(org.springframework.integration.kafka.support.ConsumerConfiguration.class);
final BeanDefinitionBuilder consumerMetadataBuilder =
BeanDefinitionBuilder.genericBeanDefinition(ConsumerMetadata.class);
BeanDefinitionBuilder.genericBeanDefinition(org.springframework.integration.kafka.support.ConsumerMetadata.class);
IntegrationNamespaceUtils.setValueIfAttributeDefined(consumerMetadataBuilder, consumerConfiguration,
"group-id");
@@ -105,7 +98,7 @@ public class KafkaConsumerContextParser extends AbstractSingleBeanDefinitionPars
if (topicFilter != null) {
BeanDefinition topicFilterConfigurationBeanDefinition =
BeanDefinitionBuilder.genericBeanDefinition(TopicFilterConfiguration.class)
BeanDefinitionBuilder.genericBeanDefinition(org.springframework.integration.kafka.support.TopicFilterConfiguration.class)
.addConstructorArgValue(topicFilter.getAttribute("pattern"))
.addConstructorArgValue(topicFilter.getAttribute("streams"))
.addConstructorArgValue(topicFilter.getAttribute("exclude"))
@@ -122,7 +115,7 @@ public class KafkaConsumerContextParser extends AbstractSingleBeanDefinitionPars
final String consumerPropertiesBean = parentElem.getAttribute("consumer-properties");
final BeanDefinitionBuilder consumerConfigFactoryBuilder =
BeanDefinitionBuilder.genericBeanDefinition(ConsumerConfigFactoryBean.class);
BeanDefinitionBuilder.genericBeanDefinition(org.springframework.integration.kafka.support.ConsumerConfigFactoryBean.class);
consumerConfigFactoryBuilder.addConstructorArgValue(consumerMetadataBeanDefintiion);
if (StringUtils.hasText(zookeeperConnectBean)) {
@@ -137,14 +130,14 @@ public class KafkaConsumerContextParser extends AbstractSingleBeanDefinitionPars
consumerConfigFactoryBuilder.getBeanDefinition();
BeanDefinitionBuilder consumerConnectionProviderBuilder =
BeanDefinitionBuilder.genericBeanDefinition(ConsumerConnectionProvider.class);
BeanDefinitionBuilder.genericBeanDefinition(org.springframework.integration.kafka.support.ConsumerConnectionProvider.class);
consumerConnectionProviderBuilder.addConstructorArgValue(consumerConfigFactoryBuilderBeanDefinition);
AbstractBeanDefinition consumerConnectionProviderBuilderBeanDefinition =
consumerConnectionProviderBuilder.getBeanDefinition();
BeanDefinitionBuilder messageLeftOverBeanDefinitionBuilder =
BeanDefinitionBuilder.genericBeanDefinition(MessageLeftOverTracker.class);
BeanDefinitionBuilder.genericBeanDefinition(org.springframework.integration.kafka.support.MessageLeftOverTracker.class);
AbstractBeanDefinition messageLeftOverBeanDefinition =
messageLeftOverBeanDefinitionBuilder.getBeanDefinition();

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2014 the original author or authors.
* Copyright 2002-2015 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
@@ -42,6 +42,7 @@ import org.springframework.util.xml.DomUtils;
* @author Soby Chacko
* @author Ilayaperumal Gopinathan
* @author Gary Russell
* @author Artem Bilan
* @since 0.5
*/
public class KafkaProducerContextParser extends AbstractSimpleBeanDefinitionParser {
@@ -120,7 +121,7 @@ public class KafkaProducerContextParser extends AbstractSimpleBeanDefinitionPars
if (StringUtils.hasText(producerConfiguration.getAttribute("partitioner"))) {
if (log.isWarnEnabled()) {
log.warn("'partitioner' is a deprecated option. Use the 'kafka_partitionId' message header or " +
"the partition argument in the send() or convertAndSend() methods");
"the 'partition-id' (or 'partition-id-expression') attribute.");
}
}
IntegrationNamespaceUtils.setValueIfAttributeDefined(producerMetadataBuilder, producerConfiguration,
@@ -131,6 +132,8 @@ public class KafkaProducerContextParser extends AbstractSimpleBeanDefinitionPars
"sync");
IntegrationNamespaceUtils.setValueIfAttributeDefined(producerMetadataBuilder, producerConfiguration,
"send-timeout");
IntegrationNamespaceUtils.setValueIfAttributeDefined(producerMetadataBuilder, producerConfiguration,
"charset");
AbstractBeanDefinition producerMetadataBeanDefinition = producerMetadataBuilder.getBeanDefinition();

View File

@@ -24,6 +24,7 @@ import kafka.utils.Utils;
*
* This class is for internal use only and therefore is at default access level
*/
@Deprecated
class DefaultPartitioner implements Partitioner {
/**
* Uses the key to calculate a partition bucket id for routing

View File

@@ -16,6 +16,8 @@
package org.springframework.integration.kafka.support;
import java.util.HashSet;
import java.util.Set;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
@@ -25,7 +27,10 @@ import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.clients.producer.RecordMetadata;
import org.apache.kafka.common.KafkaException;
import org.springframework.core.convert.ConversionService;
import org.springframework.core.convert.TypeDescriptor;
import org.springframework.core.convert.converter.GenericConverter;
import org.springframework.core.convert.support.GenericConversionService;
import org.springframework.core.serializer.support.SerializingConverter;
import org.springframework.util.Assert;
@@ -57,6 +62,7 @@ public class ProducerConfiguration<K, V> {
this.producerMetadata = producerMetadata;
this.producer = producer;
GenericConversionService genericConversionService = new GenericConversionService();
genericConversionService.addConverter(new StringBytesConverter());
genericConversionService.addConverter(Object.class, byte[].class, new SerializingConverter());
this.conversionService = genericConversionService;
}
@@ -170,4 +176,26 @@ public class ProducerConfiguration<K, V> {
'}';
}
private class StringBytesConverter implements GenericConverter {
@Override
public Set<ConvertiblePair> getConvertibleTypes() {
Set<ConvertiblePair> convertiblePairs = new HashSet<ConvertiblePair>();
convertiblePairs.add(new ConvertiblePair(String.class, byte[].class));
convertiblePairs.add(new ConvertiblePair(byte[].class, String.class));
return convertiblePairs;
}
@Override
public Object convert(Object source, TypeDescriptor sourceType, TypeDescriptor targetType) {
if (source instanceof String) {
return ((String) source).getBytes(ProducerConfiguration.this.producerMetadata.getCharset());
}
else {
return new String((byte[]) source, ProducerConfiguration.this.producerMetadata.getCharset());
}
}
}
}

View File

@@ -16,6 +16,8 @@
package org.springframework.integration.kafka.support;
import java.nio.charset.Charset;
import org.apache.kafka.common.serialization.Serializer;
import org.springframework.util.Assert;
@@ -53,6 +55,8 @@ public class ProducerMetadata<K,V> {
private boolean sync = false;
private Charset charset = Charset.forName("UTF8");
public ProducerMetadata(final String topic, Class<K> keyClassType, Class<V> valueClassType,
Serializer<K> keySerializer, Serializer<V> valueSerializer) {
Assert.notNull(topic, "Topic cannot be null");
@@ -129,6 +133,20 @@ public class ProducerMetadata<K,V> {
this.sync = sync;
}
public Charset getCharset() {
return charset;
}
/**
* The character encoding to preform {@code String <-> byte[]} conversion
* instead of general (de)serialization.
* @param charset the charset encoding to use.
* @since 1.3
*/
public void setCharset(Charset charset) {
this.charset = charset;
}
@Override
public String toString() {
return "ProducerMetadata{" +
@@ -142,6 +160,7 @@ public class ProducerMetadata<K,V> {
", batchBytes=" + this.batchBytes +
", sync=" + this.sync +
", sendTimeout=" + this.sendTimeout +
", charset=" + this.charset +
'}';
}
@@ -150,4 +169,5 @@ public class ProducerMetadata<K,V> {
gzip,
snappy
}
}

View File

@@ -216,7 +216,10 @@
<xsd:annotation>
<xsd:appinfo>
<xsd:documentation>
[DEPRECATED]
Custom Kafka key partitioner.
Deprecated in favor of 'partition-id' ('partition-id-expression')
on the 'outbound-channel-adapter'.
</xsd:documentation>
<tool:annotation kind="ref">
<tool:expected-type type="kafka.producer.Partitioner"/>
@@ -245,6 +248,14 @@
</xsd:documentation>
</xsd:annotation>
</xsd:attribute>
<xsd:attribute name="charset" use="optional" type="xsd:string">
<xsd:annotation>
<xsd:documentation>
The character encoding to preform String to/from byte[] conversion
instead of general (de)serialization. Defaults to UTF8.
</xsd:documentation>
</xsd:annotation>
</xsd:attribute>
</xsd:complexType>
</xsd:element>
</xsd:choice>

View File

@@ -44,9 +44,9 @@
key-serializer="stringSerializer"
value-serializer="stringSerializer"
batch-bytes="9876"
partitioner="partitioner"
conversion-service="conversionService"
producer-listener="producerListener"
charset="cp1251"
compression-type="none"/>
</int-kafka:producer-configurations>
</int-kafka:producer-context>
@@ -59,7 +59,6 @@
<bean id="stringSerializer" class="org.apache.kafka.common.serialization.StringSerializer"/>
<bean id="partitioner" class="org.springframework.integration.kafka.support.DefaultPartitioner"/>
<bean id="conversionService" class="org.springframework.integration.kafka.config.xml.KafkaProducerContextParserTests.StubConversionService"/>
</beans>

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2013-2014 the original author or authors.
* Copyright 2013-2015 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
@@ -22,8 +22,6 @@ import static org.junit.Assert.assertSame;
import java.util.Map;
import kafka.producer.Partitioner;
import kafka.serializer.Encoder;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.common.serialization.Serializer;
import org.junit.Assert;
@@ -38,7 +36,6 @@ import org.springframework.core.convert.ConversionService;
import org.springframework.core.convert.TypeDescriptor;
import org.springframework.integration.kafka.rule.KafkaEmbedded;
import org.springframework.integration.kafka.rule.KafkaRule;
import org.springframework.integration.kafka.rule.KafkaRunning;
import org.springframework.integration.kafka.support.KafkaProducerContext;
import org.springframework.integration.kafka.support.ProducerConfiguration;
import org.springframework.integration.kafka.support.ProducerListener;
@@ -48,6 +45,8 @@ import org.springframework.integration.test.util.TestUtils;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
import kafka.serializer.Encoder;
/**
* @author Soby Chacko
* @author Gary Russell
@@ -104,9 +103,6 @@ public class KafkaProducerContextParserTests {
assertSame(stringSerializer, producerConfigurationTest2.getProducerMetadata().getKeySerializer());
assertSame(stringSerializer, producerConfigurationTest2.getProducerMetadata().getValueSerializer());
final Partitioner partitioner = appContext.getBean("partitioner", Partitioner.class);
assertSame(partitioner, producerConfigurationTest2.getProducerMetadata().getPartitioner());
final ConversionService conversionService = appContext.getBean("conversionService", ConversionService.class);
ConversionService configuredConversionService = (ConversionService) directFieldAccessor2.getPropertyValue("conversionService");
assertSame(conversionService, configuredConversionService);
@@ -119,6 +115,8 @@ public class KafkaProducerContextParserTests {
assertFalse(TestUtils.getPropertyValue(producerContext, "autoStartup", Boolean.class));
assertEquals(123, TestUtils.getPropertyValue(producerContext, "phase"));
assertEquals("windows-1251", producerConfigurationTest2.getProducerMetadata().getCharset().name());
}
public static class StubConversionService implements ConversionService {

View File

@@ -35,6 +35,8 @@ import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import org.apache.kafka.clients.producer.RecordMetadata;
import org.apache.kafka.common.serialization.ByteArraySerializer;
import org.apache.kafka.common.serialization.StringSerializer;
import org.junit.After;
import org.junit.Rule;
import org.junit.Test;
@@ -133,7 +135,7 @@ public class OutboundTests {
KafkaProducerMessageHandler handler =
new KafkaProducerMessageHandler(producerContext);
handler.handleMessage(MessageBuilder.withPayload("foo" + suffix)
handler.handleMessage(MessageBuilder.withPayload(("foo" + suffix).getBytes())
.setHeader(KafkaHeaders.MESSAGE_KEY, "3")
.setHeader(KafkaHeaders.TOPIC, TOPIC)
.build());
@@ -298,15 +300,28 @@ public class OutboundTests {
kafkaMessageListenerContainer.start();
KafkaProducerContext producerContext = createProducerContext();
KafkaProducerMessageHandler handler = new KafkaProducerMessageHandler(producerContext);
KafkaProducerContext kafkaProducerContext = new KafkaProducerContext();
ProducerMetadata<String, byte[]> producerMetadata =
new ProducerMetadata<>(TOPIC, String.class, byte[].class,
new StringSerializer(), new ByteArraySerializer());
Properties props = new Properties();
ProducerFactoryBean<String, byte[]> producer =
new ProducerFactoryBean<>(producerMetadata, kafkaRule.getBrokersAsString(), props);
ProducerConfiguration<String, byte[]> config =
new ProducerConfiguration<>(producerMetadata, producer.getObject());
Map<String, ProducerConfiguration<?, ?>> producerConfigurationMap =
Collections.<String, ProducerConfiguration<?, ?>>singletonMap(TOPIC, config);
kafkaProducerContext.setProducerConfigurations(producerConfigurationMap);
KafkaProducerMessageHandler handler = new KafkaProducerMessageHandler(kafkaProducerContext);
handler.setBeanFactory(mock(BeanFactory.class));
handler.afterPropertiesSet();
handler.handleMessage(MessageBuilder.withPayload("fooTopic1" + suffix).build());
producerContext.stop();
kafkaProducerContext.stop();
latch.await(1000, TimeUnit.MILLISECONDS);
assertThat(latch.getCount(), equalTo(0L));

View File

@@ -21,6 +21,8 @@ import java.io.ObjectInputStream;
import java.util.ArrayList;
import kafka.producer.Partitioner;
import kafka.serializer.DefaultEncoder;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.PartitionInfo;
@@ -30,7 +32,6 @@ import org.junit.Assert;
import org.junit.Test;
import org.mockito.ArgumentCaptor;
import org.mockito.Mockito;
import org.springframework.core.convert.ConversionFailedException;
import org.springframework.integration.kafka.serializer.avro.AvroReflectDatumBackedKafkaEncoder;
import org.springframework.integration.kafka.test.utils.NonSerializableTestKey;
@@ -41,11 +42,10 @@ import org.springframework.integration.kafka.util.EncoderAdaptingSerializer;
import org.springframework.messaging.Message;
import org.springframework.messaging.support.GenericMessage;
import kafka.serializer.DefaultEncoder;
/**
* @author Soby Chacko
* @author Artem Bilan
* @author Marius Bogoevici
* @since 0.5
*/
public class ProducerConfigurationTests {
@@ -154,9 +154,7 @@ public class ProducerConfigurationTests {
final byte[] keyBytes = capturedKeyMessage.key();
final ByteArrayInputStream keyInputStream = new ByteArrayInputStream(keyBytes);
final ObjectInputStream keyObjectInputStream = new ObjectInputStream(keyInputStream);
final Object keyObj = keyObjectInputStream.readObject();
String keyObj = new String(keyBytes);
Assert.assertEquals("key", keyObj);
Assert.assertEquals(capturedKeyMessage.value(), tp);
@@ -195,11 +193,8 @@ public class ProducerConfigurationTests {
final byte[] payloadBytes = capturedKeyMessage.value();
final ByteArrayInputStream payloadBis = new ByteArrayInputStream(payloadBytes);
final ObjectInputStream payloadOis = new ObjectInputStream(payloadBis);
final Object payloadObj = payloadOis.readObject();
Assert.assertEquals("test message", payloadObj);
String payload = new String(payloadBytes);
Assert.assertEquals("test message", payload);
Assert.assertEquals(capturedKeyMessage.topic(), "test");
}
@@ -228,19 +223,13 @@ public class ProducerConfigurationTests {
final ProducerRecord<byte[], byte[]> capturedKeyMessage = argument.getValue();
final byte[] keyBytes = capturedKeyMessage.key();
final ByteArrayInputStream keyBis = new ByteArrayInputStream(keyBytes);
final ObjectInputStream keyOis = new ObjectInputStream(keyBis);
final Object keyObj = keyOis.readObject();
Assert.assertEquals("key", keyObj);
String key = new String(keyBytes);
Assert.assertEquals("key", key);
final byte[] payloadBytes = capturedKeyMessage.value();
final ByteArrayInputStream payloadBis = new ByteArrayInputStream(payloadBytes);
final ObjectInputStream payloadOis = new ObjectInputStream(payloadBis);
final Object payloadObj = payloadOis.readObject();
Assert.assertEquals("test message", payloadObj);
String payload = new String(payloadBytes);
Assert.assertEquals("test message", payload);
Assert.assertEquals(capturedKeyMessage.topic(), "test");
}