Skip to content

将 MQTT 数据导入 BigQuery

BigQuery 是一个完全托管的企业级数据仓库,适用于大规模、基于 SQL 的分析和报表。EMQX Cloud 可通过规则引擎和 BigQuery Sink 将 MQTT 数据流式传输到 BigQuery,实现物联网数据的实时提取、处理与分析。

本页介绍如何在 EMQX Cloud 中创建 BigQuery 数据集成。示例将 MQTT 主题 test/a 中的消息写入 BigQuery 数据表。

工作原理

BigQuery 数据集成使用 EMQX Cloud 规则引擎选择和转换 MQTT 消息,再通过 BigQuery Sink 将规则输出发送到 BigQuery。

数据流如下:

text
MQTT 客户端 -> EMQX Cloud -> 规则 -> BigQuery Sink -> BigQuery 数据表
  1. MQTT 客户端向 test/a 等主题发布遥测数据或事件数据。
  2. 规则匹配来自该主题的消息,并选择需要写入的字段。
  3. BigQuery Sink 将所选字段写入配置的 BigQuery 数据集和数据表。
  4. 您可以在 BigQuery 中查询该数据表,对导入的 MQTT 数据进行分析。

准备工作

前置准备

开始操作前,请确保您已了解:

配置网络访问

BigQuery 连接器通过 HTTPS 连接到 BigQuery。请根据部署类型配置网络:

  • 对于弹性专有版部署,如果部署需要通过公网访问 Google Cloud 服务,请启用 NAT 网关
  • 对于 BYOC 部署,请确保部署所在的 VPC 可以访问 BigQuery。如果需要通过公网访问,请在云服务商控制台中配置 NAT 网关。

在 GCP 中创建服务账户密钥

要允许 EMQX Cloud 向 BigQuery 写入数据,请在 Google Cloud 中创建服务账户并生成 JSON 格式的密钥。

  1. 在 Google Cloud 项目中创建一个服务账户

  2. 为服务账户授予向目标 BigQuery 数据集和数据表写入数据所需的权限。例如,为目标数据集授予 BigQuery Data Editor 角色,或根据您的安全策略授予同等的读写权限。

  3. 打开服务账户详情页面,点击密钥页签,创建一个 JSON 格式的新密钥。

    TIP

    请妥善保管下载的服务账户密钥。创建 EMQX Cloud BigQuery 连接器时需要提供此密钥。

在 BigQuery 中创建数据集和数据表

在 EMQX Cloud 中配置 BigQuery Sink 之前,请先在 Google Cloud 中创建目标数据集和数据表。

  1. 在 Google Cloud 控制台中,进入 BigQuery -> Studio

  2. Explorer 窗格中创建数据集。例如,创建名为 emq_test_dataset 的数据集。

  3. 在该数据集中创建数据表。例如,创建名为 bigquery_integration_test 的数据表。

  4. 定义数据表 Schema。本教程使用以下 Schema:

    text
    clientid:string,payload:bytes,topic:string,publish_received_at:timestamp
  5. 确认已创建的服务账户具有目标数据表的写入权限。

  6. 可选:运行以下查询,检查是否可以访问该数据表。请将项目、数据集和数据表名称替换为您实际使用的名称。

    sql
    SELECT * FROM `my_project.emq_test_dataset.bigquery_integration_test` LIMIT 1000

创建 BigQuery 连接器

创建规则前,请先创建 BigQuery 连接器,以连接 EMQX Cloud 和 BigQuery。

  1. 在 EMQX Cloud 控制台中进入您的部署。

  2. 在左侧导航菜单中点击数据集成

  3. 如果这是您创建的第一个连接器,请在数据持久化分类下选择 BigQuery。如果已经存在连接器,请点击新建连接器,然后选择 BigQuery

  4. 新建连接器页面中配置以下字段:

    • 连接器名称:使用系统自动生成的名称,或输入名称。
    • GCP 服务账户凭证:粘贴在在 GCP 中创建服务账户密钥中创建的服务账户密钥的完整 JSON 内容,或点击选择文件导入 JSON 文件。
    • 其他设置使用默认值,或根据业务需求进行配置。
  5. 点击测试以验证连接。如果 BigQuery 服务可访问且服务账户凭证有效,系统将返回成功提示。

  6. 点击新建完成连接器设置。现在您可以创建规则并添加 BigQuery Sink 动作。

创建规则

创建规则,以选择需要写入 BigQuery 的 MQTT 消息字段。

  1. 规则区域点击新建规则,或点击连接器旁的操作图标。

  2. SQL 编辑器中输入以下 SQL:

    sql
    SELECT
      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 规则。

  3. 点击下一步,为规则添加动作。

添加 BigQuery Sink

新建动作 (Sink) 页面中配置 BigQuery Sink,将规则输出写入 BigQuery。

  1. 配置动作:

    • 连接器:选择已创建的 BigQuery 连接器。
    • 动作类型:值为 BigQuery
    • 动作名称:使用系统自动生成的名称,或输入名称。
    • 数据集:输入 BigQuery 数据集名称,例如 emq_test_dataset
    • 数据表:输入 BigQuery 数据表名称,例如 bigquery_integration_test
  2. 可选:如果需要提高消息传递失败时的可靠性,请配置备用动作。当主 BigQuery Sink 无法处理消息时,将触发备用动作。

  3. 除非需要调整连接或缓冲行为,否则请保持高级设置的默认值。有关详细信息,请参阅高级设置

  4. 点击确认,创建规则和动作。

  5. 成功创建规则弹窗中点击返回规则列表,完成规则创建。

测试规则

使用 MQTTX 或其他 MQTT 客户端向 test/a 主题发布测试消息。

  1. 向 EMQX Cloud 发布以下消息:

    bash
    mqttx pub -i c_emqx -t test/a -m '{ "msg": "hello" }'
  2. 在 EMQX Cloud 控制台中进入规则列表,点击规则 ID,查看规则和动作统计。规则应新增一条流入消息,BigQuery Sink 应新增一条流出消息。

  3. 在 Google Cloud 控制台中进入 BigQuery -> Studio,打开目标数据表并运行以下查询。请将项目、数据集和数据表名称替换为您实际使用的名称。

    sql
    SELECT *
    FROM `my_project.emq_test_dataset.bigquery_integration_test`
    ORDER BY publish_received_at DESC
    LIMIT 10

    您应能看到写入目标数据表的消息。

    BigQuery 查询结果

如果查询返回了测试消息,表示集成正常工作:

text
MQTT -> 规则 -> BigQuery Sink -> BigQuery 数据表

高级设置

配置 BigQuery Sink 时,可以展开高级设置,根据需要调整以下参数。

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