在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
相关产品推荐
相关产品推荐

