#283: Added basic Kafka testcase that executes a workflow

This commit is contained in:
frikky
2021-05-29 18:42:41 +02:00
parent 35336663d7
commit 5423685321
4 changed files with 88 additions and 0 deletions
+4
View File
@@ -0,0 +1,4 @@
```
curl -sSL https://raw.githubusercontent.com/bitnami/bitnami-docker-kafka/master/docker-compose.yml > docker-compose.yml
docker-compose up -d
```
@@ -0,0 +1,28 @@
version: "2"
services:
zookeeper:
image: docker.io/bitnami/zookeeper:3
ports:
- "2181:2181"
volumes:
- "zookeeper_data:/bitnami"
environment:
- ALLOW_ANONYMOUS_LOGIN=yes
kafka:
image: docker.io/bitnami/kafka:2
ports:
- "9092:9092"
volumes:
- "kafka_data:/bitnami"
environment:
- KAFKA_CFG_ZOOKEEPER_CONNECT=zookeeper:2181
- ALLOW_PLAINTEXT_LISTENER=yes
depends_on:
- zookeeper
volumes:
zookeeper_data:
driver: local
kafka_data:
driver: local
+54
View File
@@ -0,0 +1,54 @@
import json
import os
import requests
from time import sleep
from kafka import KafkaProducer, KafkaConsumer
import kafka
shuffle_url = os.getenv("SHUFFLE_URL")
shuffle_apikey = os.getenv("SHUFFLE_APIKEY")
shuffle_workflow = os.getenv("SHUFFLE_WORKFLOW")
headers = {"Authorization": "Bearer %s" % shuffle_apikey}
topic = "workflow_%s" % shuffle_workflow
server = "localhost:9092"
def produce():
print("Starting producer")
producer = KafkaProducer(
bootstrap_servers=[server],
value_serializer=lambda x:
json.dumps(x).encode('utf-8')
)
print("Adding data!")
for e in range(15):
data = {"some": e, "data": "luuuuul"}
try:
ret = producer.send(topic, value=data)
print(ret.get())
except kafka.errors.KafkaTimeoutError as e:
print("Kafka error: %s" % e)
continue
def consume():
print("Starting consumer")
consumer = KafkaConsumer(
topic,
bootstrap_servers=[server],
auto_offset_reset="earliest",
enable_auto_commit=True,
value_deserializer=lambda x: json.loads(x.decode('utf-8'))
)
print("Getting data")
for message in consumer:
message = message.value
print("MSG: ", message)
ret = requests.post("%s/api/v1/%s/execute" % shuffle_url, headers=headers, data=message)
print(ret.status_code)
print(ret.text)
if __name__ == "__main__":
produce()
#consume()
@@ -0,0 +1,2 @@
kafka
kafka-python