前置条件
在腾讯云环境中使用 AutoMQ Table Topic,需要满足以下条件:-
版本要求: AutoMQ 实例版本为
1.4.1或更高版本。 - Catalog 要求: 准备可访问的 Hive Metastore 服务。你可以自行部署 Hive Metastore,也可以使用腾讯云 EMR 提供的 HMS。
-
存储要求: 提前创建 COS Bucket,并确认 AutoMQ 实例所在地域与 Bucket 地域一致。记下包含 APPID 后缀的完整 Bucket 标识、地域和 endpoint,例如
examplebucket-1250000000、ap-shanghai和https://cos.ap-shanghai.myqcloud.com。 - 权限要求: 为 AutoMQ 实例使用的 CAM Role 授予访问 COS Bucket 的权限。创建 Hive Catalog 集成时,控制台会生成所需的 Policy。
操作步骤
步骤 1:准备 Hive Metastore 和 COS
配置 EMR HMS
- 前往腾讯云 EMR 产品,创建或选择一个包含 Hive Metastore 的 EMR 集群。如果创建新集群,请选择 Hive 组件;后续查询还需要 Spark 组件。根据你的安全要求决定是否开启 Kerberos。

- 在 EMR 控制台进入集群服务 > HIVE > 配置管理,选择
hive-env.sh,点击编辑配置,添加访问对象存储所需的 Jar。以下配置适用于本文示例环境;如果 EMR 版本不同,请根据实际 Hadoop 版本调整 Jar 路径和版本。

- 在 HMS 使用的
hive-site.xml中配置 COS 的 S3 兼容访问参数。以下示例使用静态凭证访问 COS,请将YOUR_TENCENT_SECRET_ID和YOUR_TENCENT_SECRET_KEY替换为具备目标 Bucket 访问权限的实际凭证。
ap-shanghai 替换为 COS Bucket 所在地域。COS endpoint 的格式和地域列表请参考腾讯云 COS 地域和访问域名。

- 保存配置并重启 Hive Metastore,使配置生效。

准备 Kerberos 文件(可选)
如果 HMS 启用了 Kerberos,请准备一个用于访问 HMS 的客户端 principal、对应的 keytab 文件和krb5.conf 文件。以下命令需要在 EMR Master 节点上以 root 身份执行;addprinc 和 ktadd 是 kadmin.local 的交互式命令:
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 集成
- 登录 AutoMQ 控制台,进入集成页面。

-
选择创建 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 路径。

- 复制控制台根据 Warehouse 生成的 COS Policy,在腾讯云 CAM 控制台创建自定义策略,并将策略附加到 AutoMQ 实例使用的 CAM Role。至少确认该 Role 对目标 Bucket 及其对象拥有 Table Topic 所需的读、写、列举和删除权限。
- 完成授权后,创建 Hive Catalog 集成。
步骤 3:创建实例并开启 Table Topic
创建 AutoMQ 实例时,在 Table Topic 配置中执行以下操作。Table Topic 依赖 AutoMQ 内置的 Schema Registry,开启 Table Topic 时请同时确保 Schema Registry 已开启。- 开启 Table Topic。
- 选择 Hive Catalog。
- 选择步骤 2 创建的 Hive Catalog 集成。

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

步骤 5:生产消息并使用 Spark SQL 查询
- 打开 Topic 详情中的生产消息页签,输入测试消息的 Key 和 Value,然后发送消息。
- 确认 EMR 集群包含 Spark 组件;如果创建集群时未选择 Spark,请先添加 Spark 组件。登录 EMR Master 节点并切换到
hadoop用户。 - 本示例使用截图中的 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 集成文档。
- 在
spark-sql中查询表。将<namespace>和<topic>替换为创建 Topic 时使用的命名空间和 Topic 名称:
