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

在GCP上用Python操作Kafka读写GCS桶时遇NotImplementedError报错求助

问题排查与解决

报错原因

你调用blob.open()时传参逻辑错误:blob.open()的第一个参数是文件打开模式(比如'w'),但你把完整的GCS路径当成第一个参数传入,导致GCS SDK误将该路径解析为模式字符串,而它不在支持的模式列表('r', 'rb', 'rt', 'w', 'wb', 'wt')中,因此抛出NotImplementedError。

同时代码逻辑存在误区:你先创建了一个名为stock_market的Blob对象,循环里却试图用这个Blob去打开另一个路径的文件——每个Blob对象对应GCS中的一个具体文件,不能跨路径复用。

修正方案

方案1:正确使用blob.open()

每次循环根据计数生成目标Blob名称,创建对应的Blob对象后,传入正确的模式参数调用open():

from google.cloud import storage
import json

storage_client = storage.Client()
bucket = storage_client.bucket("kafka-gcptest")

for count, msg in enumerate(consumer):
    # 生成每个消息对应的Blob文件名
    blob_name = f'stock_market_{count}.json'
    blob = bucket.blob(blob_name)
    # 第一个参数传入正确的打开模式
    with blob.open('w') as file:
        json.dump(msg.value, file)

方案2:使用upload_from_string更高效

如果只是写入JSON内容,直接用upload_from_string可以省去文件对象操作,更简洁高效:

from google.cloud import storage
import json

storage_client = storage.Client()
bucket = storage_client.bucket("kafka-gcptest")

for count, msg in enumerate(consumer):
    blob_name = f'stock_market_{count}.json'
    blob = bucket.blob(blob_name)
    # 将JSON对象转为字符串后上传,同时指定内容类型
    blob.upload_from_string(json.dumps(msg.value), content_type='application/json')

内容的提问来源于stack exchange,提问作者Arpit Gupta

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 20:44:55