带Schema的Pub/Sub推送订阅无法转发消息如何解决?
解决带Avro Schema的Pub/Sub Push订阅消息解析问题
问题核心原因
当你给Pub/Sub Topic绑定Avro Schema后,默认情况下Pub/Sub会将发布的消息序列化为二进制Avro格式(从消息属性googclient_schemaencoding: BINARY可验证),而非保留原始JSON。BigQuery订阅内置了Avro Schema解析能力,能自动完成解码,而自定义Push订阅需要手动处理二进制Avro数据,直接用UTF-8解码自然会失败。
两种解决方案
方案一:正确解析二进制Avro数据
你需要用Topic绑定的Avro Schema对base64解码后的二进制流进行解析,确保Schema与Topic绑定的完全一致(包括字段名、类型、顺序,Avro二进制编码对字段顺序敏感)。
Python示例代码
import base64 import io import avro.schema from avro.io import BinaryDecoder, DatumReader # 1. 复制Topic绑定的Avro Schema定义 schema_str = """ { "type": "record", "name": "Metrics", "fields": [ {"name": "id", "type": "int"}, {"name": "monitoring_id", "type": "string"}, {"name": "timestamp_micros", "type": "long"}, {"name": "updated_rows", "type": "int"}, {"name": "deleted_rows", "type": "int"}, {"name": "inserted_rows", "type": "int"}, {"name": "bad_formatted_rows", "type": "int"} ] } """ schema = avro.schema.parse(schema_str) # 2. 处理Push订阅收到的消息数据 received_data = "AhBBQkxVVVVVVQIUAAAU" binary_data = base64.b64decode(received_data) # 3. 解码Avro二进制流 reader = DatumReader(schema) decoder = BinaryDecoder(io.BytesIO(binary_data)) decoded_msg = reader.read(decoder) print(decoded_msg) # 输出:{'id': 1, 'monitoring_id': 'ABLUUUUU', 'timestamp_micros': 1, 'updated_rows': 10, 'deleted_rows': 0, 'inserted_rows': 0, 'bad_formatted_rows': 10}
Node.js示例代码
const avro = require('avsc'); const { Buffer } = require('buffer'); // 1. 复制Topic绑定的Avro Schema定义 const schema = avro.parse({ type: 'record', name: 'Metrics', fields: [ { name: 'id', type: 'int' }, { name: 'monitoring_id', type: 'string' }, { name: 'timestamp_micros', type: 'long' }, { name: 'updated_rows', type: 'int' }, { name: 'deleted_rows', type: 'int' }, { name: 'inserted_rows', type: 'int' }, { name: 'bad_formatted_rows', type: 'int' } ] }); // 2. 处理Push订阅收到的消息数据 const receivedData = 'AhBBQkxVVVVVVQIUAAAU'; const binaryData = Buffer.from(receivedData, 'base64'); // 3. 解码Avro二进制流 const decodedMsg = schema.fromBuffer(binaryData); console.log(decodedMsg);
方案二:让Push订阅转发原始JSON数据
如果业务允许,你可以在发布消息时指定JSON编码,让Pub/Sub保留原始JSON格式,这样Push订阅收到的data就是原始JSON的base64编码,直接解码即可得到结构化内容。
用gcloud CLI发布消息示例
gcloud pubsub topics publish YOUR_TOPIC_NAME \ --message '{"id":1,"monitoring_id":"ABLUUUUU","timestamp_micros":1,"updated_rows":10,"deleted_rows":0,"inserted_rows":0,"bad_formatted_rows":10}' \ --attribute googclient_schemaencoding=JSON,googclient_schemarevisionid=404dc06d
发布后,Push订阅收到的data字段会是原始JSON的base64字符串,直接用UTF-8解码就能得到你预期的结构化数据。
内容的提问来源于stack exchange,提问作者Joao Leme
相关产品推荐
相关产品推荐

