Skip to content

Instantly share code, notes, and snippets.

@danielkec
Created June 17, 2020 18:50
Show Gist options
  • Select an option

  • Save danielkec/8a12cc970ec872c3174b4f2b7ce4a7c8 to your computer and use it in GitHub Desktop.

Select an option

Save danielkec/8a12cc970ec872c3174b4f2b7ce4a7c8 to your computer and use it in GitHub Desktop.
SE Messaging with MP config
mp.messaging:
incoming.from-kafka:
connector: helidon-kafka
topic: messaging-test-topic-1
auto.offset.reset: latest
enable.auto.commit: true
outgoing.to-kafka:
connector: helidon-kafka
topic: messaging-test-topic-1
connector:
helidon-kafka:
bootstrap.servers: localhost:9092
key.serializer: org.apache.kafka.common.serialization.StringSerializer
value.serializer: org.apache.kafka.common.serialization.StringSerializer
key.deserializer: org.apache.kafka.common.serialization.StringDeserializer
value.deserializer: org.apache.kafka.common.serialization.StringDeserializer
Channel<String> fromKafka = Channel.<String>builder()
.name("from-kafka")
.publisherConfig(KafkaConnector.configBuilder()
.groupId("example-group-" + session.getId())
.build()
)
.build();
// Prepare Kafka connector, can be used by any channel
KafkaConnector kafkaConnector = KafkaConnector.create();
Messaging messaging = Messaging.builder()
.connector(kafkaConnector)
.listener(fromKafka, payload -> {
System.out.println("Kafka says: " + payload);
// Send message received from Kafka over websocket
sendTextMessage(session, payload);
})
.build()
.start();
// Prepare channel for connecting processor -> kafka connector with specific subscriber configuration,
// channel -> connector mapping is automatic when using KafkaConnector.configBuilder()
Channel<String> toKafka = Channel.<String>create("to-kafka");
// Prepare channel for connecting emitter -> processor
Channel<String> toProcessor = Channel.create();
// Prepare Kafka connector, can be used by any channel
KafkaConnector kafkaConnector = KafkaConnector.create();
// Prepare emitter for manual publishing to channel
emitter = Emitter.create(toProcessor);
messaging = Messaging.builder()
.emitter(emitter)
// Processor connect two channels together
.processor(toProcessor, toKafka, payload -> {
// Transforming to upper-case before sending to kafka
return payload.toUpperCase();
})
.connector(kafkaConnector)
.build()
.start();
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment