Skip to main content
AutoMQ Table Topic 可以将 Kafka Topic 数据持续写入 Apache Iceberg 表。本文以腾讯云 EMR 的 Hive Metastore(HMS)和腾讯云对象存储(COS)为例,介绍如何配置集成并使用 Spark SQL 查询数据。 Table Topic 在后台批量提交 Iceberg snapshot,因此数据不会在每条消息写入后立即出现在查询结果中。测试时,请等待首个 snapshot 提交完成后再查询。

前置条件

在腾讯云环境中使用 AutoMQ Table Topic,需要满足以下条件:
  • 版本要求: AutoMQ 实例版本为 1.4.1 或更高版本。
  • Catalog 要求: 准备可访问的 Hive Metastore 服务。你可以自行部署 Hive Metastore,也可以使用腾讯云 EMR 提供的 HMS。
  • 存储要求: 提前创建 COS Bucket,并确认 AutoMQ 实例所在地域与 Bucket 地域一致。记下包含 APPID 后缀的完整 Bucket 标识、地域和 endpoint,例如 examplebucket-1250000000ap-shanghaihttps://cos.ap-shanghai.myqcloud.com
  • 权限要求: 为 AutoMQ 实例使用的 CAM Role 授予访问 COS Bucket 的权限。创建 Hive Catalog 集成时,控制台会生成所需的 Policy。
本文使用 EMR HMS、Kerberos 和 Spark SQL 作为示例。Kerberos 不是必选项;如果你的 HMS 未启用 Kerberos,可以跳过相关步骤并在控制台选择无认证

操作步骤

步骤 1:准备 Hive Metastore 和 COS

配置 EMR HMS

  1. 前往腾讯云 EMR 产品,创建或选择一个包含 Hive Metastore 的 EMR 集群。如果创建新集群,请选择 Hive 组件;后续查询还需要 Spark 组件。根据你的安全要求决定是否开启 Kerberos
腾讯云 EMR 产品购买页面,开启 Kerberos 身份认证
  1. 在 EMR 控制台进入集群服务 > HIVE > 配置管理,选择 hive-env.sh,点击编辑配置,添加访问对象存储所需的 Jar。以下配置适用于本文示例环境;如果 EMR 版本不同,请根据实际 Hadoop 版本调整 Jar 路径和版本。
配置 hive-env.sh 文件,添加访问对象存储所需的 Jar 依赖 不要同时加载多套不同版本的 Hadoop AWS 或 AWS SDK Jar。
  1. 在 HMS 使用的 hive-site.xml 中配置 COS 的 S3 兼容访问参数。以下示例使用静态凭证访问 COS,请将 YOUR_TENCENT_SECRET_IDYOUR_TENCENT_SECRET_KEY 替换为具备目标 Bucket 访问权限的实际凭证。
ap-shanghai 替换为 COS Bucket 所在地域。COS endpoint 的格式和地域列表请参考腾讯云 COS 地域和访问域名 配置 hive-site.xml 文件,添加对象存储访问相关配置
  1. 保存配置并重启 Hive Metastore,使配置生效。
Hive Metastore 服务重启配置页面

准备 Kerberos 文件(可选)

如果 HMS 启用了 Kerberos,请准备一个用于访问 HMS 的客户端 principal、对应的 keytab 文件和 krb5.conf 文件。以下命令需要在 EMR Master 节点上以 root 身份执行;addprincktaddkadmin.local 的交互式命令:
复制 Kerberos 配置文件,并在后续步骤上传这两个文件:
如果当前 EMR 集群使用了旧的 des3-cbc-sha1 加密类型,只有在该集群的 Kerberos 配置要求时,才在 /root/krb5.conf[libdefaults] 段添加 allow_weak_crypto = true。不要在不需要时启用弱加密。 在 Hive 服务配置中获取以下值:
  • hive.metastore.uris:Hive Metastore 接入点,格式为 thrift://<host>:<port>。它填入 AutoMQ 的 Hive Metastore Endpoint,不是 User Principal。
  • hive.metastore.kerberos.principal:Hive Metastore 服务端 principal,填入 AutoMQ 的 Kerberos Principal
  • automq-user@YOUR_REALM:上面创建的客户端 principal,填入 AutoMQ 的 Kerberos User Principal

步骤 2:创建 Hive Catalog 集成

  1. 登录 AutoMQ 控制台,进入集成页面。
AutoMQ 控制台集成菜单
  1. 选择创建 Hive Catalog 集成,填写以下信息:
    • 名称: 填写便于识别的集成名称。
    • 部署配置: 选择 AutoMQ 实例所属的部署配置。该配置必须与后续创建的实例一致。
    • Hive Metastore Endpoint: 填写步骤 1 中获取的 hive.metastore.uris,例如 thrift://hms.example.com:7004
    • 鉴权模式: 根据 HMS 配置选择无认证Kerberos
    • Keytab 文件: 使用 Kerberos 时,上传客户端 keytab 文件。
    • krb5.conf 文件: 使用 Kerberos 时,上传 krb5.conf 文件。
    • Kerberos User Principal: 使用 Kerberos 时,填写客户端 principal,例如 automq-user@YOUR_REALM
    • Kerberos Principal: 使用 Kerberos 时,填写 hive.metastore.kerberos.principal 的值,即 HMS 服务端 principal。
    • Warehouse: 填写包含 APPID 后缀的完整 COS Bucket 标识,不要填写 s3a:// 前缀。AutoMQ 会使用 <bucket>/iceberg 作为 Iceberg Warehouse 路径。
创建 Hive Catalog 集成表单,包含 Hive Metastore 接入点、鉴权模式、Kerberos 参数和 Warehouse 字段
  1. 复制控制台根据 Warehouse 生成的 COS Policy,在腾讯云 CAM 控制台创建自定义策略,并将策略附加到 AutoMQ 实例使用的 CAM Role。至少确认该 Role 对目标 Bucket 及其对象拥有 Table Topic 所需的读、写、列举和删除权限。
  2. 完成授权后,创建 Hive Catalog 集成。

步骤 3:创建实例并开启 Table Topic

创建 AutoMQ 实例时,在 Table Topic 配置中执行以下操作。Table Topic 依赖 AutoMQ 内置的 Schema Registry,开启 Table Topic 时请同时确保 Schema Registry 已开启。
  1. 开启 Table Topic
  2. 选择 Hive Catalog
  3. 选择步骤 2 创建的 Hive Catalog 集成。
Table Topic 开启后,并非所有 Topic 都会自动写入 Iceberg 表。你仍需在 Topic 粒度开启流转表。已选择的 Catalog 配置提交后不能修改;如需更换 Catalog,需要创建新实例。 AutoMQ 实例创建页面,已开启 Table Topic 功能

步骤 4:创建 Topic 并配置流转表

  1. 进入实例的 Topics 页面,点击创建 Topic
  2. 配置 Topic 的基础参数:
    • Topic 名称: 填写 Topic 名称,该名称也会作为 Iceberg 表名。
  3. 开启 Table Topic 转换,并配置以下参数:
    • 命名空间: Iceberg Catalog 中的 Database 名称,用于隔离不同的 Iceberg 表。
    • Schema 约束类型:
      • Schema: 消息必须遵守已注册的 Schema。需要配置AutoMQ 内置的 Schema Registry并注册消息 Schema;Table Topic 使用 Schema 的字段创建和填充 Iceberg 表。
      • Schemaless: 消息没有明确的 Schema,Table Topic 将消息 Key 和 Value 作为整体字段写入 Iceberg 表。
  4. 点击确定创建 Topic。
Topic 创建表单,已开启 Table Topic 转换,显示命名空间和 Schema 约束类型字段

步骤 5:生产消息并使用 Spark SQL 查询

  1. 打开 Topic 详情中的生产消息页签,输入测试消息的 Key 和 Value,然后发送消息。
  2. 确认 EMR 集群包含 Spark 组件;如果创建集群时未选择 Spark,请先添加 Spark 组件。登录 EMR Master 节点并切换到 hadoop 用户。
  3. 本示例使用截图中的 Spark 3.2.2、Scala 2.12 和 Iceberg 0.13.1。若 EMR 使用其他版本,请同步调整以下变量,确保 iceberg-spark-runtime 的 Spark/Scala 后缀与实际环境匹配。
spark.sql.catalog.hive.s3.* 配置用于 Iceberg 的 S3FileIO;命令中的 spark.hadoop.fs.s3a.* 配置保留原 EMR Spark/Hadoop 的 S3A 访问参考。两套配置分别作用于不同的访问实现。有关 S3FileIO 和依赖配置,请参考 Apache Iceberg AWS 集成文档
  1. spark-sql 中查询表。将 <namespace><topic> 替换为创建 Topic 时使用的命名空间和 Topic 名称:
Table Topic 提交 Iceberg snapshot 后,Spark SQL 可以查询对应的表记录,并返回查询结果行数。你也可以使用其他支持 Hive Catalog 和 Iceberg 的查询引擎,但需要分别配置 Catalog、COS 权限和网络连通性。 Spark SQL 查询结果,显示 AutoMQ 将 Kafka 消息转换为 Iceberg 表数据记录