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

带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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 13:26:29