能否在Google Colab单Notebook中运行Kafka生产者与消费者?
在单个Google Colab Notebook中同时运行Kafka生产者和消费者
核心思路
完全可以实现,在Colab里通过后台线程或子进程就能让生产者和消费者并行执行,不用分开终端。下面是具体落地步骤:
步骤1:先在Colab里搭好Kafka环境
Colab默认不带Kafka,先运行以下命令安装并启动依赖服务:
# 安装Java(Kafka必须依赖) !apt-get install openjdk-8-jdk-headless -qq > /dev/null # 下载解压Kafka包 !wget -q https://downloads.apache.org/kafka/3.6.1/kafka_2.13-3.6.1.tgz !tar xf kafka_2.13-3.6.1.tgz # 后台启动Zookeeper !./kafka_2.13-3.6.1/bin/zookeeper-server-start.sh -daemon ./kafka_2.13-3.6.1/config/zookeeper.properties # 后台启动Kafka Broker !./kafka_2.13-3.6.1/bin/kafka-server-start.sh -daemon ./kafka_2.13-3.6.1/config/server.properties # 等10秒让服务启动,然后创建测试主题 !sleep 10 !./kafka_2.13-3.6.1/bin/kafka-topics.sh --create --topic test_topic --bootstrap-server localhost:9092 --partitions 1 --replication-factor 1
步骤2:用Python线程实现并行运行
以confluent-kafka库为例(也适配你仓库里的Spark Kafka逻辑),先装依赖:
!pip install confluent-kafka
然后写代码,用threading把生产者和消费者拆到不同线程:
from confluent_kafka import Producer, Consumer, KafkaError import threading import time import os # 生产者配置 producer_conf = {'bootstrap.servers': 'localhost:9092'} producer = Producer(producer_conf) # 消费者配置 consumer_conf = { 'bootstrap.servers': 'localhost:9092', 'group.id': 'test_group', 'auto.offset.reset': 'earliest' } consumer = Consumer(consumer_conf) consumer.subscribe(['test_topic']) # 生产者任务:循环发10条测试消息 def run_producer(): for i in range(10): msg_content = f"测试消息{i}" producer.produce('test_topic', value=msg_content.encode('utf-8')) producer.flush() print(f"已发送: {msg_content}") time.sleep(1) # 消费者任务:持续消费消息,直到生产者结束 def run_consumer(): while True: msg = consumer.poll(1.0) if msg is None: continue if msg.error(): if msg.error().code() == KafkaError._PARTITION_EOF: continue print(msg.error()) break print(f"已接收: {msg.value().decode('utf-8')}") consumer.close() # 启动两个线程 producer_thread = threading.Thread(target=run_producer) consumer_thread = threading.Thread(target=run_consumer) producer_thread.start() consumer_thread.start() # 等生产者发完消息,再停掉消费者 producer_thread.join() time.sleep(2) os._exit(0)
步骤3:适配你仓库里的Spark Kafka代码
如果要沿用仓库中的Spark实现,逻辑类似:
- 把仓库里的生产者代码封装成独立函数,消费者代码也封装成函数
- 分别放到不同线程启动,注意共用同一个SparkSession,避免重复创建上下文
额外注意事项
- Colab会话有超时限制,长时间运行可能断开,测试时尽量用短周期任务
- 如果遇到连接失败,先查Kafka进程是否正常:
!ps aux | grep kafka - 消费者线程可以加个全局标志位来优雅停止,不用强制退出主线程
内容的提问来源于stack exchange,提问作者Dorsa Arezooji
相关产品推荐
相关产品推荐

