Вот мой докер, созданный для Zookeeper и Kafka:
Код: Выделить всё
zookeeper:
image: confluentinc/cp-zookeeper:7.6.0
hostname: zookeeper
container_name: zookeeper
ports:
- "32181:32181"
environment:
ZOOKEEPER_CLIENT_PORT: 32181
ZOOKEEPER_SERVER_ID: 1
ZOOKEEPER_TICK_TIME: 2000
ZOOKEEPER_SYNC_LIMIT: 2
healthcheck:
test: ['CMD', 'bash', '-c', "echo 'ruok' | nc localhost 32181"]
interval: 10s
timeout: 5s
retries: 5
networks:
- airflow-network
kafka:
image: confluentinc/cp-kafka:7.6.0
hostname: kafka
container_name: kafka
depends_on:
zookeeper:
condition: service_healthy
ports:
- "9092:9092"
- "29092:29092"
expose:
- 9092
environment:
KAFKA_ZOOKEEPER_CONNECT: zookeeper:32181
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092,PLAINTEXT_HOST://localhost:29092
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
KAFKA_BROKER_ID: 1
KAFKA_INIT_OFFSET_CHANGES: "true"
KAFKA_NUM_PARTITIONS: 1
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
healthcheck:
test: [ "CMD", "bash", "-c", 'nc -z localhost 9092']
interval: 10s
timeout: 5s
retries: 5
networks:
- airflow-network
Код: Выделить всё
from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime
from confluent_kafka import Producer
from confluent_kafka import Consumer
import json
import requests
import socket
default_args = {
'owner': 'airflow',
'start_date': datetime(2024, 6, 4, 13, 00),
}
def create_kafka_producer():
config = {'bootstrap.servers': 'kafka:9092',
'client.id': socket.gethostname()}
producer = Producer(config)
return producer
def create_kafka_consumer():
conf = {'bootstrap.servers': 'host1:9092,host2:9092',
'group.id': 'foo',
'auto.offset.reset': 'smallest'}
consumer = Consumer(conf)
return consumer
def get_data():
res = requests.get('https://randomuser.me/api')
res = res.json()
res = res['results'][0]
return res
def format_data(res):
data = {}
location = res['location']
data['first_name'] = res['name']['first']
data['last_name'] = res['name']['last']
data['gender'] = res['gender']
data['address'] = f"{str(location['street']['number'])} " \
f"{location['street']['name']}, " \
f"{location['city']}, {location['state']}, " \
f"{location['country']}"
data['post_code'] = location['postcode']
data['email'] = res['email']
data['username'] = res['login']['username']
data['dob'] = res['dob']['date']
data['registered_date'] = res['registered']['date']
data['phone'] = res['phone']
data['picture'] = res['picture']['medium']
return data
def stream_data(res):
res = get_data()
res = format_data(res)
producer = create_kafka_producer()
producer.produce('user_data',
json.dumps(res).encode('utf-8'),
callback=delivery_callback)
def delivery_callback(err, msg):
if err:
print('ERROR: Message failed delivery: {}'.format(err))
else:
print("Produced event to topic")
with DAG('user_automation',
default_args=default_args,
schedule_interval='@daily',
catchup=False) as dag:
streaming_task = PythonOperator(
task_id='stream_data_from_api',
python_callable=stream_data
)
streaming_task
Код: Выделить всё
docker exec kafka kafka-console-consumer --bootstrap-server kafka:9092 --topic user_data --from-beginningПодробнее здесь: https://stackoverflow.com/questions/783 ... ed-nothing