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

使用GCP DataFlow迁移数据时Pubsub消息超100字节限制问题求助

问题分析与解决:Pubsub消息长度超过限制

报错原因

你当前使用TableRow.toPrettyString()方法生成消息payload,这个方法会输出带换行、缩进的格式化字符串,导致单条消息的字节体积远大于实际需要,触发了Pubsub的消息大小限制(报错提示的100字节上限,大概率是测试环境的特殊配额限制)。

解决方法

1. 改用紧凑JSON序列化

替换toPrettyString()为紧凑的JSON序列化方式,消除多余的空格和换行,大幅缩小消息体积:

// 引入Gson依赖(如果项目中没有的话)
import com.google.gson.Gson

// 在DoFn中修改序列化逻辑
private val gson = Gson()
val data = gson.toJson(row).toByteArray()

2. 只保留必要字段

如果业务不需要传递TableRow的全部字段,仅提取核心字段构造消息,进一步减少payload大小:

val essentialFields = mapOf(
    "user_id" to row.get("user_id"),
    "order_amount" to row.get("order_amount")
)
val data = gson.toJson(essentialFields).toByteArray()

3. 确认Pubsub配额配置

若上述调整后仍报错,检查目标Pubsub Topic是否设置了自定义的消息大小限制,或DataFlow作业中是否有相关限流配置(默认Pubsub单条消息上限为10MB,100字节属于特殊限制场景)。

修改后的完整代码

import com.google.gson.Gson
import org.apache.beam.sdk.Pipeline
import org.apache.beam.sdk.io.BigQueryIO
import org.apache.beam.sdk.io.gcp.pubsub.PubsubIO
import org.apache.beam.sdk.transforms.DoFn
import org.apache.beam.sdk.transforms.ParDo
import com.google.cloud.bigquery.TableRow
import com.google.pubsub.v1.PubsubMessage

protected fun executeDataflow(pipeline: Pipeline, query: String, destinationTopic: String) {
    pipeline
            .apply("BigQueryRead", BigQueryIO.readTableRows()
                    .fromQuery(query)
                    .usingStandardSql()
                    .withQueryLocation("EU")
            )
            .apply("convert to message",
                ParDo.of(object : DoFn<TableRow, PubsubMessage>() {
                    private val gson = Gson()
                    
                    @ProcessElement
                    fun processElement(context: ProcessContext) {
                        val row = context.element()
                        val data = gson.toJson(row).toByteArray()
                        val pubsubMessage = PubsubMessage(data, null)

                        context.output(pubsubMessage)
                    }
                })
            )
            .apply("Pubsub", PubsubIO.writeMessages().to(destinationTopic))

    pipeline.run()
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 21:00:15