diff --git a/functions/extensions/kafka/README.md b/functions/extensions/kafka/README.md new file mode 100644 index 00000000..dceb22a3 --- /dev/null +++ b/functions/extensions/kafka/README.md @@ -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 +``` diff --git a/functions/extensions/kafka/docker-compose.yml b/functions/extensions/kafka/docker-compose.yml new file mode 100644 index 00000000..9a4af72a --- /dev/null +++ b/functions/extensions/kafka/docker-compose.yml @@ -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 diff --git a/functions/extensions/kafka/kafka_local.py b/functions/extensions/kafka/kafka_local.py new file mode 100644 index 00000000..a9e16eb9 --- /dev/null +++ b/functions/extensions/kafka/kafka_local.py @@ -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() diff --git a/functions/extensions/kafka/requirements.txt b/functions/extensions/kafka/requirements.txt new file mode 100644 index 00000000..652ae858 --- /dev/null +++ b/functions/extensions/kafka/requirements.txt @@ -0,0 +1,2 @@ +kafka +kafka-python