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 MB、32 MB 或 64 MB(随部署连接数变化)
请求模式可选择同步或异步请求模式。在异步模式下,写入 BigQuery 不会阻塞 MQTT 消息发布过程。Async
批量大小指定单个批次中发送到 BigQuery 的最大记录数。如果设置为 1,则逐条发送记录。1000
请求飞行队列窗口指定与 BigQuery 通信时允许同时存在的最大飞行请求数。当请求模式为 Async 时,如果需要严格按顺序处理,请将此值设置为 1。100