Потребитель консоли Kafka ничего не вывелPython

Программы на Python
Anonymous
Потребитель консоли Kafka ничего не вывел

Сообщение Anonymous »

Я новичок в Кафке. Развернут Kafka и Airflow внутри докер-контейнера. Кажется, мой датчик воздушного потока работает правильно. Я пытался отправлять сообщения в Kafka, но когда я использую консольного потребителя Kafka, он ничего не вывел.
Вот мой докер, созданный для 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
Вот мой DAG в воздушном потоке

Код: Выделить всё

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
Вот моя потребительская команда консоли Kafka, которую я выполнил.

Код: Выделить всё

docker exec kafka kafka-console-consumer --bootstrap-server kafka:9092 --topic user_data --from-beginning
Я ожидал, что эта команда выведет сообщения, но она ничего не вывела.

Подробнее здесь: https://stackoverflow.com/questions/783 ... ed-nothing

Вернуться в «Python»