You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

能否在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.21 09:40:16