-
Notifications
You must be signed in to change notification settings - Fork 217
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
#586 add configuration for kafka consumer, add migration notes for ch…
…anged kafka configuration Signed-off-by: Johannes Schneider <johannes.schneider@bosch.io>
- Loading branch information
Showing
22 changed files
with
282 additions
and
250 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
45 changes: 0 additions & 45 deletions
45
...ava/org/eclipse/ditto/connectivity/service/messaging/kafka/ConsumerPropertiesFactory.java
This file was deleted.
Oops, something went wrong.
79 changes: 0 additions & 79 deletions
79
...org/eclipse/ditto/connectivity/service/messaging/kafka/DefaultKafkaConnectionFactory.java
This file was deleted.
Oops, something went wrong.
54 changes: 54 additions & 0 deletions
54
...a/org/eclipse/ditto/connectivity/service/messaging/kafka/DefaultKafkaProducerFactory.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,54 @@ | ||
/* | ||
* Copyright (c) 2017 Contributors to the Eclipse Foundation | ||
* | ||
* See the NOTICE file(s) distributed with this work for additional | ||
* information regarding copyright ownership. | ||
* | ||
* This program and the accompanying materials are made available under the | ||
* terms of the Eclipse Public License 2.0 which is available at | ||
* http://www.eclipse.org/legal/epl-2.0 | ||
* | ||
* SPDX-License-Identifier: EPL-2.0 | ||
*/ | ||
package org.eclipse.ditto.connectivity.service.messaging.kafka; | ||
|
||
import java.util.Map; | ||
|
||
import org.apache.kafka.clients.producer.KafkaProducer; | ||
import org.apache.kafka.clients.producer.Producer; | ||
import org.apache.kafka.common.serialization.Serializer; | ||
import org.apache.kafka.common.serialization.StringSerializer; | ||
import org.eclipse.ditto.connectivity.model.Connection; | ||
import org.eclipse.ditto.connectivity.service.config.KafkaConfig; | ||
|
||
/** | ||
* Default implementation of {@link org.eclipse.ditto.connectivity.service.messaging.kafka.KafkaProducerFactory}. | ||
*/ | ||
final class DefaultKafkaProducerFactory implements KafkaProducerFactory { | ||
|
||
private static final Serializer<String> KEY_SERIALIZER = new StringSerializer(); | ||
private static final Serializer<String> VALUE_SERIALIZER = KEY_SERIALIZER; | ||
|
||
private final Map<String, Object> producerProperties; | ||
|
||
private DefaultKafkaProducerFactory(final Map<String, Object> producerProperties) { | ||
this.producerProperties = producerProperties; | ||
} | ||
|
||
/** | ||
* Returns an instance of the default Kafka connection factory. | ||
* | ||
* @param propertiesFactory a factory to create kafka client configuration properties. | ||
* @return an Kafka connection factory. | ||
*/ | ||
static DefaultKafkaProducerFactory getInstance(final PropertiesFactory propertiesFactory) { | ||
final Map<String, Object> producerProperties = propertiesFactory.getProducerProperties(); | ||
return new DefaultKafkaProducerFactory(producerProperties); | ||
} | ||
|
||
@Override | ||
public Producer<String, String> newProducer() { | ||
return new KafkaProducer<>(producerProperties, KEY_SERIALIZER, VALUE_SERIALIZER); | ||
} | ||
|
||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.