Move message header mode as a generic property
- Both the producer and consumer properties have `HeaderMode` - Handle the case of embeddedHeaders and raw for both the Sending/ReceivingHandlers in Redis binder This resolves #408 Move message values extraction to superclass
This commit is contained in:
committed by
Marius Bogoevici
parent
cd81f24451
commit
502d275b41
@@ -33,8 +33,6 @@ public class KafkaConsumerProperties {
|
||||
|
||||
private KafkaMessageChannelBinder.StartOffset startOffset = null;
|
||||
|
||||
private KafkaMessageChannelBinder.Mode mode = KafkaMessageChannelBinder.Mode.embeddedHeaders;
|
||||
|
||||
public void setMinPartitionCount(int minPartitionCount) {
|
||||
this.minPartitionCount = minPartitionCount;
|
||||
}
|
||||
@@ -52,14 +50,6 @@ public class KafkaConsumerProperties {
|
||||
this.autoCommitOffset = autoCommitOffset;
|
||||
}
|
||||
|
||||
public KafkaMessageChannelBinder.Mode getMode() {
|
||||
return mode;
|
||||
}
|
||||
|
||||
public void setMode(KafkaMessageChannelBinder.Mode mode) {
|
||||
this.mode = mode;
|
||||
}
|
||||
|
||||
public boolean isResetOffsets() {
|
||||
return resetOffsets;
|
||||
}
|
||||
|
||||
@@ -29,16 +29,9 @@ import java.util.Properties;
|
||||
import java.util.UUID;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
|
||||
import kafka.admin.AdminUtils;
|
||||
import kafka.api.OffsetRequest;
|
||||
import kafka.serializer.Decoder;
|
||||
import kafka.serializer.DefaultDecoder;
|
||||
import kafka.utils.ZKStringSerializer$;
|
||||
import kafka.utils.ZkUtils;
|
||||
import org.I0Itec.zkclient.ZkClient;
|
||||
import org.apache.kafka.clients.producer.ProducerConfig;
|
||||
import org.apache.kafka.common.serialization.ByteArraySerializer;
|
||||
import scala.collection.Seq;
|
||||
|
||||
import org.springframework.beans.factory.DisposableBean;
|
||||
import org.springframework.beans.factory.config.ConfigurableListableBeanFactory;
|
||||
@@ -48,10 +41,10 @@ import org.springframework.cloud.stream.binder.BinderException;
|
||||
import org.springframework.cloud.stream.binder.BinderHeaders;
|
||||
import org.springframework.cloud.stream.binder.Binding;
|
||||
import org.springframework.cloud.stream.binder.DefaultBinding;
|
||||
import org.springframework.cloud.stream.binder.EmbeddedHeadersMessageConverter;
|
||||
import org.springframework.cloud.stream.binder.ExtendedConsumerProperties;
|
||||
import org.springframework.cloud.stream.binder.ExtendedProducerProperties;
|
||||
import org.springframework.cloud.stream.binder.ExtendedPropertiesBinder;
|
||||
import org.springframework.cloud.stream.binder.HeaderMode;
|
||||
import org.springframework.cloud.stream.binder.MessageValues;
|
||||
import org.springframework.cloud.stream.binder.PartitionHandler;
|
||||
import org.springframework.http.MediaType;
|
||||
@@ -90,6 +83,14 @@ import org.springframework.util.Assert;
|
||||
import org.springframework.util.CollectionUtils;
|
||||
import org.springframework.util.StringUtils;
|
||||
|
||||
import kafka.admin.AdminUtils;
|
||||
import kafka.api.OffsetRequest;
|
||||
import kafka.serializer.Decoder;
|
||||
import kafka.serializer.DefaultDecoder;
|
||||
import kafka.utils.ZKStringSerializer$;
|
||||
import kafka.utils.ZkUtils;
|
||||
import scala.collection.Seq;
|
||||
|
||||
/**
|
||||
* A {@link Binder} that uses Kafka as the underlying middleware.
|
||||
*
|
||||
@@ -109,9 +110,6 @@ public class KafkaMessageChannelBinder extends AbstractBinder<MessageChannel, Ex
|
||||
|
||||
private final Map<String, Collection<Partition>> topicsInUse = new HashMap<>();
|
||||
|
||||
private final EmbeddedHeadersMessageConverter embeddedHeadersMessageConverter = new
|
||||
EmbeddedHeadersMessageConverter();
|
||||
|
||||
private final ZookeeperConnect zookeeperConnect;
|
||||
|
||||
private final String brokers;
|
||||
@@ -577,17 +575,8 @@ public class KafkaMessageChannelBinder extends AbstractBinder<MessageChannel, Ex
|
||||
@Override
|
||||
@SuppressWarnings("unchecked")
|
||||
protected Object handleRequestMessage(Message<?> requestMessage) {
|
||||
if (Mode.embeddedHeaders.equals(consumerProperties.getExtension().getMode())) {
|
||||
MessageValues messageValues;
|
||||
try {
|
||||
messageValues = embeddedHeadersMessageConverter.extractHeaders((Message<byte[]>) requestMessage,
|
||||
true);
|
||||
}
|
||||
catch (Exception e) {
|
||||
logger.error(EmbeddedHeadersMessageConverter.decodeExceptionMessage(requestMessage), e);
|
||||
messageValues = new MessageValues(requestMessage);
|
||||
}
|
||||
messageValues = deserializePayloadIfNecessary(messageValues);
|
||||
if (HeaderMode.embeddedHeaders.equals(consumerProperties.getHeaderMode())) {
|
||||
MessageValues messageValues = extractMessageValues(requestMessage);
|
||||
return MessageBuilder.createMessage(messageValues.getPayload(), new KafkaBinderHeaders(
|
||||
messageValues));
|
||||
}
|
||||
@@ -648,14 +637,13 @@ public class KafkaMessageChannelBinder extends AbstractBinder<MessageChannel, Ex
|
||||
else {
|
||||
targetPartition = roundRobin() % numberOfKafkaPartitions;
|
||||
}
|
||||
|
||||
if (Mode.embeddedHeaders.equals(producerProperties.getExtension().getMode())) {
|
||||
if (HeaderMode.embeddedHeaders.equals(producerProperties.getHeaderMode())) {
|
||||
MessageValues transformed = serializePayloadIfNecessary(message);
|
||||
byte[] messageToSend = embeddedHeadersMessageConverter.embedHeaders(transformed,
|
||||
KafkaMessageChannelBinder.this.headersToMap);
|
||||
producerConfiguration.send(topicName, targetPartition, null, messageToSend);
|
||||
}
|
||||
else if (Mode.raw.equals(producerProperties.getExtension().getMode())) {
|
||||
else if (HeaderMode.raw.equals(producerProperties.getHeaderMode())) {
|
||||
Object contentType = message.getHeaders().get(MessageHeaders.CONTENT_TYPE);
|
||||
if (contentType != null
|
||||
&& !contentType.equals(MediaType.APPLICATION_OCTET_STREAM_VALUE)) {
|
||||
@@ -682,11 +670,6 @@ public class KafkaMessageChannelBinder extends AbstractBinder<MessageChannel, Ex
|
||||
|
||||
}
|
||||
|
||||
public enum Mode {
|
||||
raw,
|
||||
embeddedHeaders
|
||||
}
|
||||
|
||||
public enum StartOffset {
|
||||
earliest(OffsetRequest.EarliestTime()),
|
||||
latest(OffsetRequest.LatestTime());
|
||||
|
||||
@@ -31,8 +31,6 @@ public class KafkaProducerProperties {
|
||||
|
||||
private boolean sync = false;
|
||||
|
||||
private KafkaMessageChannelBinder.Mode mode = KafkaMessageChannelBinder.Mode.embeddedHeaders;
|
||||
|
||||
private int batchTimeout = 0;
|
||||
|
||||
public int getBufferSize() {
|
||||
@@ -60,15 +58,6 @@ public class KafkaProducerProperties {
|
||||
this.sync = sync;
|
||||
}
|
||||
|
||||
@NotNull
|
||||
public KafkaMessageChannelBinder.Mode getMode() {
|
||||
return mode;
|
||||
}
|
||||
|
||||
public void setMode(KafkaMessageChannelBinder.Mode mode) {
|
||||
this.mode = mode;
|
||||
}
|
||||
|
||||
public int getBatchTimeout() {
|
||||
return batchTimeout;
|
||||
}
|
||||
|
||||
@@ -17,8 +17,6 @@
|
||||
package org.springframework.cloud.stream.binder.kafka.config;
|
||||
|
||||
import org.springframework.boot.context.properties.ConfigurationProperties;
|
||||
import org.springframework.cloud.stream.binder.kafka.KafkaMessageChannelBinder;
|
||||
import org.springframework.cloud.stream.binder.kafka.KafkaMessageChannelBinder.Mode;
|
||||
import org.springframework.util.StringUtils;
|
||||
|
||||
/**
|
||||
@@ -39,8 +37,6 @@ class KafkaBinderConfigurationProperties {
|
||||
|
||||
private String[] headers = new String[] {};
|
||||
|
||||
private KafkaMessageChannelBinder.Mode mode = Mode.embeddedHeaders;
|
||||
|
||||
private int offsetUpdateTimeWindow = 10000;
|
||||
|
||||
private int offsetUpdateCount = 0;
|
||||
@@ -81,10 +77,6 @@ class KafkaBinderConfigurationProperties {
|
||||
return headers;
|
||||
}
|
||||
|
||||
public Mode getMode() {
|
||||
return this.mode;
|
||||
}
|
||||
|
||||
public int getOffsetUpdateTimeWindow() {
|
||||
return this.offsetUpdateTimeWindow;
|
||||
}
|
||||
@@ -118,10 +110,6 @@ class KafkaBinderConfigurationProperties {
|
||||
this.headers = headers;
|
||||
}
|
||||
|
||||
public void setMode(KafkaMessageChannelBinder.Mode mode) {
|
||||
this.mode = mode;
|
||||
}
|
||||
|
||||
public void setOffsetUpdateTimeWindow(int offsetUpdateTimeWindow) {
|
||||
this.offsetUpdateTimeWindow = offsetUpdateTimeWindow;
|
||||
}
|
||||
|
||||
@@ -75,6 +75,9 @@ public abstract class AbstractBinder<T, C extends ConsumerProperties, P extends
|
||||
|
||||
private final StringConvertingContentTypeResolver contentTypeResolver = new StringConvertingContentTypeResolver();
|
||||
|
||||
protected final EmbeddedHeadersMessageConverter embeddedHeadersMessageConverter = new
|
||||
EmbeddedHeadersMessageConverter();
|
||||
|
||||
protected volatile EvaluationContext evaluationContext;
|
||||
|
||||
protected volatile PartitionSelectorStrategy partitionSelector;
|
||||
@@ -139,6 +142,26 @@ public abstract class AbstractBinder<T, C extends ConsumerProperties, P extends
|
||||
onInit();
|
||||
}
|
||||
|
||||
/**
|
||||
* Extract the message values from the the received message when the received message is embedded with
|
||||
* header values. Once extracted, deserialize the payload if necessary.
|
||||
*
|
||||
* @param receivedMessage the received message
|
||||
* @return extracted message values
|
||||
*/
|
||||
public MessageValues extractMessageValues(Message<?> receivedMessage) {
|
||||
MessageValues messageValues;
|
||||
try {
|
||||
messageValues = embeddedHeadersMessageConverter.extractHeaders((Message<byte[]>) receivedMessage,
|
||||
true);
|
||||
}
|
||||
catch (Exception e) {
|
||||
logger.error(EmbeddedHeadersMessageConverter.decodeExceptionMessage(receivedMessage), e);
|
||||
messageValues = new MessageValues(receivedMessage);
|
||||
}
|
||||
return deserializePayloadIfNecessary(messageValues);
|
||||
}
|
||||
|
||||
/**
|
||||
* Subclasses may implement this method to perform any necessary initialization.
|
||||
* It will be invoked from {@link #afterPropertiesSet()} which is itself {@code final}.
|
||||
|
||||
@@ -47,6 +47,8 @@ public class ConsumerProperties {
|
||||
|
||||
private double backOffMultiplier = 2.0;
|
||||
|
||||
private HeaderMode headerMode = HeaderMode.embeddedHeaders;
|
||||
|
||||
@Min(value = 1, message = "Concurrency should be greater than zero.")
|
||||
public int getConcurrency() {
|
||||
return concurrency;
|
||||
@@ -118,5 +120,12 @@ public class ConsumerProperties {
|
||||
return backOffMultiplier;
|
||||
}
|
||||
|
||||
public HeaderMode getHeaderMode() {
|
||||
return this.headerMode;
|
||||
}
|
||||
|
||||
public void setHeaderMode(HeaderMode headerMode) {
|
||||
this.headerMode = headerMode;
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,24 @@
|
||||
/*
|
||||
* Copyright 2016 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.
|
||||
* You may obtain a copy of the License at
|
||||
*
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
*
|
||||
* Unless required by applicable law or agreed to in writing, software
|
||||
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
package org.springframework.cloud.stream.binder;
|
||||
|
||||
/**
|
||||
* @author Marius Bogoevici
|
||||
*/
|
||||
public enum HeaderMode {
|
||||
raw,
|
||||
embeddedHeaders
|
||||
}
|
||||
@@ -48,6 +48,8 @@ public class ProducerProperties {
|
||||
|
||||
private String[] requiredGroups = new String[] {};
|
||||
|
||||
private HeaderMode headerMode = HeaderMode.embeddedHeaders;
|
||||
|
||||
public Expression getPartitionKeyExpression() {
|
||||
return partitionKeyExpression;
|
||||
}
|
||||
@@ -111,4 +113,12 @@ public class ProducerProperties {
|
||||
return (this.partitionSelectorClass == null) || (this.partitionSelectorExpression == null);
|
||||
}
|
||||
|
||||
public HeaderMode getHeaderMode() {
|
||||
return this.headerMode;
|
||||
}
|
||||
|
||||
public void setHeaderMode(HeaderMode headerMode) {
|
||||
this.headerMode = headerMode;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user