deployment.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
name: quickstart-mp-deployment
labels:
app: quickstart-mp
spec:
selector:| public class ExampleBean { | |
| private final SubmissionPublisher<String> publisher = new SubmissionPublisher<>(); | |
| public void sendMessage(String message) { | |
| publisher.submit(message); | |
| } | |
| @Outgoing("to-kafka") | |
| public Publisher<String> preparePublisher() { |
| Channel<String> fromKafka = Channel.create("from-kafka"); | |
| KafkaConnector kafkaConnector = KafkaConnector.create(); | |
| Messaging.builder() | |
| .connector(kafkaConnector) | |
| .listener(fromKafka, payload -> { | |
| System.out.println("Kafka says: " + payload); | |
| }) | |
| .build() | |
| .start(); |
| Channel<String> toKafka = Channel.create("to-kafka"); | |
| KafkaConnector kafkaConnector = KafkaConnector.create(); | |
| Emitter<String> emitter = Emitter.create(toKafka); | |
| Messaging.builder() | |
| .emitter(emitter) | |
| .connector(kafkaConnector) | |
| .build() | |
| .start(); |
| mp.messaging: | |
| incoming.from-kafka: | |
| connector: helidon-kafka | |
| topic: messaging-test-topic-1 | |
| enable.auto.commit: false | |
| ack.timeout.millis: 10000 | |
| group.id: example-group-1 |
| Channel<String> fromKafka = Channel.create("from-kafka"); | |
| KafkaConnector kafkaConnector = KafkaConnector.create(); | |
| Messaging messaging = Messaging.builder() | |
| .connector(kafkaConnector) | |
| .subscriber(fromKafka, ReactiveStreams.<Message<String>>builder() | |
| //Apply back-pressure, flatMapCompletionStage requests one by one | |
| .flatMapCompletionStage(message -> { | |
| return CompletableFuture.runAsync(() -> { | |
| //Do something lengthy |
| Channel<String> fromKafka = Channel.<String>builder() | |
| .publisherConfig(KafkaConnector.configBuilder() | |
| .bootstrapServers("localhost:9092") | |
| .groupId("example-group-1") | |
| .topic("messaging-test-topic-1") | |
| .autoOffsetReset(KafkaConfigBuilder.AutoOffsetReset.LATEST) | |
| .enableAutoCommit(false) | |
| .keyDeserializer(StringDeserializer.class) | |
| .valueDeserializer(StringDeserializer.class) | |
| .build() |
| <dependency> | |
| <groupId>io.helidon.microprofile.messaging</groupId> | |
| <artifactId>helidon-microprofile-messaging</artifactId> | |
| </dependency> | |
| <dependency> | |
| <groupId>io.helidon.messaging.kafka</groupId> | |
| <artifactId>helidon-messaging-kafka</artifactId> | |
| </dependency> |
| ############################################# | |
| # PROVIDE THESE VALUES NEEDED FOR WORKSHOP... | |
| ############################################# | |
| export GRAALVM_HOME=~/graalvm-ce-java11-20.1.0 | |
| export JAEGER_QUERY_ADDRESS=http://130.61.207.39:80 | |
| #for example export JAEGER_QUERY_ADDRESS=http://123.123.123.123:8080 | |
| export OCI_REGION=eu-frankfurt-1 | |
| # for example export OCI_REGION=us-phoenix-1 |
| kubectl get pods --all-namespaces | |
| kubectl get svc --all-namespaces | |
| kubectl logs --tail=100 -l app=order -n msdataworkshop | |
| kubectl logs --tail=100 -l app=inventory -n msdataworkshop | |
| kubectl logs -f -l app=order -n msdataworkshop | |
| kubectl logs -f -l app=inventory -n msdataworkshop | |
| kubectl cluster-info dump | |
| kubectl delete pod,service order |
deployment.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
name: quickstart-mp-deployment
labels:
app: quickstart-mp
spec:
selector: