将 MQTT 数据导入 BigQuery
BigQuery 是一个完全托管的企业级数据仓库,适用于大规模、基于 SQL 的分析和报表。EMQX Cloud 可通过规则引擎和 BigQuery Sink 将 MQTT 数据流式传输到 BigQuery,实现物联网数据的实时提取、处理与分析。
本页介绍如何在 EMQX Cloud 中创建 BigQuery 数据集成。示例将 MQTT 主题 test/a 中的消息写入 BigQuery 数据表。
工作原理
BigQuery 数据集成使用 EMQX Cloud 规则引擎选择和转换 MQTT 消息,再通过 BigQuery Sink 将规则输出发送到 BigQuery。
数据流如下:
MQTT 客户端 -> EMQX Cloud -> 规则 -> BigQuery Sink -> BigQuery 数据表- MQTT 客户端向
test/a等主题发布遥测数据或事件数据。 - 规则匹配来自该主题的消息,并选择需要写入的字段。
- BigQuery Sink 将所选字段写入配置的 BigQuery 数据集和数据表。
- 您可以在 BigQuery 中查询该数据表,对导入的 MQTT 数据进行分析。
准备工作
前置准备
开始操作前,请确保您已了解:
配置网络访问
BigQuery 连接器通过 HTTPS 连接到 BigQuery。请根据部署类型配置网络:
- 对于弹性专有版部署,如果部署需要通过公网访问 Google Cloud 服务,请启用 NAT 网关。
- 对于 BYOC 部署,请确保部署所在的 VPC 可以访问 BigQuery。如果需要通过公网访问,请在云服务商控制台中配置 NAT 网关。
在 GCP 中创建服务账户密钥
要允许 EMQX Cloud 向 BigQuery 写入数据,请在 Google Cloud 中创建服务账户并生成 JSON 格式的密钥。
在 Google Cloud 项目中创建一个服务账户。
为服务账户授予向目标 BigQuery 数据集和数据表写入数据所需的权限。例如,为目标数据集授予 BigQuery Data Editor 角色,或根据您的安全策略授予同等的读写权限。
打开服务账户详情页面,点击密钥页签,创建一个 JSON 格式的新密钥。
TIP
请妥善保管下载的服务账户密钥。创建 EMQX Cloud BigQuery 连接器时需要提供此密钥。
在 BigQuery 中创建数据集和数据表
在 EMQX Cloud 中配置 BigQuery Sink 之前,请先在 Google Cloud 中创建目标数据集和数据表。
在 Google Cloud 控制台中,进入 BigQuery -> Studio。
在 Explorer 窗格中创建数据集。例如,创建名为
emq_test_dataset的数据集。在该数据集中创建数据表。例如,创建名为
bigquery_integration_test的数据表。定义数据表 Schema。本教程使用以下 Schema:
textclientid:string,payload:bytes,topic:string,publish_received_at:timestamp确认已创建的服务账户具有目标数据表的写入权限。
可选:运行以下查询,检查是否可以访问该数据表。请将项目、数据集和数据表名称替换为您实际使用的名称。
sqlSELECT * FROM `my_project.emq_test_dataset.bigquery_integration_test` LIMIT 1000
创建 BigQuery 连接器
创建规则前,请先创建 BigQuery 连接器,以连接 EMQX Cloud 和 BigQuery。
在 EMQX Cloud 控制台中进入您的部署。
在左侧导航菜单中点击数据集成。
如果这是您创建的第一个连接器,请在数据持久化分类下选择 BigQuery。如果已经存在连接器,请点击新建连接器,然后选择 BigQuery。
在新建连接器页面中配置以下字段:
- 连接器名称:使用系统自动生成的名称,或输入名称。
- GCP 服务账户凭证:粘贴在在 GCP 中创建服务账户密钥中创建的服务账户密钥的完整 JSON 内容,或点击选择文件导入 JSON 文件。
- 其他设置使用默认值,或根据业务需求进行配置。
点击测试以验证连接。如果 BigQuery 服务可访问且服务账户凭证有效,系统将返回成功提示。
点击新建完成连接器设置。现在您可以创建规则并添加 BigQuery Sink 动作。
创建规则
创建规则,以选择需要写入 BigQuery 的 MQTT 消息字段。
在规则区域点击新建规则,或点击连接器旁的操作图标。
在 SQL 编辑器中输入以下 SQL:
sqlSELECT clientid, topic, base64_encode(payload) AS payload, timestamp/1000 AS publish_received_at FROM "test/a"此规则监听发送到
test/a主题的消息,并仅选择与 BigQuery 数据表 Schema 匹配的字段。TIP
BigQuery 不接受未知字段。如果您自定义 SQL 或数据表 Schema,请确保所选字段名称与目标数据表中的列名一致。
TIP
如果您是初次使用规则 SQL,可以点击试运行了解并测试 SQL 规则。
点击下一步,为规则添加动作。
添加 BigQuery Sink
在新建动作 (Sink) 页面中配置 BigQuery Sink,将规则输出写入 BigQuery。
配置动作:
- 连接器:选择已创建的 BigQuery 连接器。
- 动作类型:值为 BigQuery。
- 动作名称:使用系统自动生成的名称,或输入名称。
- 数据集:输入 BigQuery 数据集名称,例如
emq_test_dataset。 - 数据表:输入 BigQuery 数据表名称,例如
bigquery_integration_test。
可选:如果需要提高消息传递失败时的可靠性,请配置备用动作。当主 BigQuery Sink 无法处理消息时,将触发备用动作。
除非需要调整连接或缓冲行为,否则请保持高级设置的默认值。有关详细信息,请参阅高级设置。
点击确认,创建规则和动作。
在成功创建规则弹窗中点击返回规则列表,完成规则创建。
测试规则
使用 MQTTX 或其他 MQTT 客户端向 test/a 主题发布测试消息。
向 EMQX Cloud 发布以下消息:
bashmqttx pub -i c_emqx -t test/a -m '{ "msg": "hello" }'在 EMQX Cloud 控制台中进入规则列表,点击规则 ID,查看规则和动作统计。规则应新增一条流入消息,BigQuery Sink 应新增一条流出消息。
在 Google Cloud 控制台中进入 BigQuery -> Studio,打开目标数据表并运行以下查询。请将项目、数据集和数据表名称替换为您实际使用的名称。
sqlSELECT * FROM `my_project.emq_test_dataset.bigquery_integration_test` ORDER BY publish_received_at DESC LIMIT 10您应能看到写入目标数据表的消息。

如果查询返回了测试消息,表示集成正常工作:
MQTT -> 规则 -> BigQuery Sink -> BigQuery 数据表高级设置
配置 BigQuery Sink 时,可以展开高级设置,根据需要调整以下参数。
| 字段名称 | 描述 | 默认值 |
|---|---|---|
| 缓冲池大小 | 指定用于管理 EMQX Cloud 与 BigQuery 之间数据流的缓冲 Worker 进程数量。这些 Worker 会在将数据发送到 BigQuery 之前临时存储并处理数据。 | 16 |
| 请求超期 | 指定请求进入缓冲区后保持有效的最长时间。如果请求在缓冲区中停留的时间超过此值,或已发送但未及时收到 BigQuery 的响应,则请求过期。 | 45 秒 |
| 健康检查间隔 | 指定 Sink 自动检查与 BigQuery 连接的时间间隔。 | 15 秒 |
| 健康检查间隔抖动 | 在基础健康检查间隔上增加随机延迟,以降低多个节点同时发起健康检查的概率。 | 0 毫秒 |
| 健康检查超时 | 指定连接器健康检查的超时时间。 | 5 秒 |
| 最大缓存队列大小 | 指定 BigQuery Sink 中每个缓冲 Worker 进程最多可缓冲的字节数。 | 16 MB、32 MB 或 64 MB(随部署连接数变化) |
| 请求模式 | 可选择同步或异步请求模式。在异步模式下,写入 BigQuery 不会阻塞 MQTT 消息发布过程。 | Async |
| 批量大小 | 指定单个批次中发送到 BigQuery 的最大记录数。如果设置为 1,则逐条发送记录。 | 1000 |
| 请求飞行队列窗口 | 指定与 BigQuery 通信时允许同时存在的最大飞行请求数。当请求模式为 Async 时,如果需要严格按顺序处理,请将此值设置为 1。 | 100 |