Skip to content

Instantly share code, notes, and snippets.

View danielkec's full-sized avatar
🚀

Daniel Kec danielkec

🚀
View GitHub Profile
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
@danielkec
danielkec / kubectl.md
Last active October 2, 2020 20:31
K8s health probes cheatsheet

deployment.yaml

apiVersion: apps/v1
kind: Deployment
metadata:
  name: quickstart-mp-deployment
  labels:
    app: quickstart-mp
spec:
 selector: