钩子

Storm

qianmoQqianmoQ· 更新于 2026-09-27· 阅读 6 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

Apache Atlas 的 Apache Storm 钩子

简介

Apache Storm 是一个分布式实时计算系统。Storm 让可靠地处理无界数据流变得简单,它在实时处理领域所起的作用,就如同 Hadoop 在批处理领域所起的作用一样。其处理过程本质上是由节点组成的 DAG,称为 topology(拓扑)。

Apache Atlas 是一个元数据存储库,支持端到端的数据血缘、搜索以及业务分类关联。

本次集成的目标是将运行中的拓扑元数据,连同其底层的数据源、目标、衍生过程以及任何可用的业务上下文一并推送,以便 Atlas 能够捕获该拓扑的血缘关系。

此过程包含以下两个部分,详见下文:

  • 用于表示 Storm 中各种概念的数据模型
  • 用于更新 Atlas 中元数据的 Storm Atlas 钩子

Storm 数据模型

数据模型在 Atlas 中以类型(Type)的形式表示。它包含拓扑图中各类节点的描述,例如 spout、bolt 以及相应的生产者和消费者类型。

Atlas 中新增了以下类型。

  • storm_topology - 表示粗粒度的拓扑。storm_topology 派生自 Atlas 的 Process 类型,因此可用于向 Atlas 提供血缘信息。
  • 新增了以下数据集 - kafka_topic、jms_topic、hbase_table、hdfs_data_set。这些都派生自 Atlas 的 Dataset 类型,因此构成血缘图的端点。
  • storm_spout - 具有输出的数据生产者,通常输出到 Kafka、JMS
  • storm_bolt - 具有输入和输出的数据消费者,通常输入输出到 Hive、HBase、HDFS 等

如果 Storm Atlas 钩子发现 Atlas 服务器尚不了解其依赖的模型(例如 Hive 数据模型),它会自动注册这些模型。

每种类型的数据模型描述位于类定义 org.apache.atlas.storm.model.StormDataModel 中。

Storm Atlas 钩子

当新的拓扑在 Storm 中成功注册时,会通知 Atlas。Storm 在用于提交 storm 拓扑的 Storm 客户端上提供了一个钩子 backtype.storm.ISubmitterHook。

Storm Atlas 钩子会拦截该钩子执行后的操作,从拓扑中提取元数据,并使用所定义的类型更新 Atlas。Atlas 在 org.apache.atlas.storm.hook.StormAtlasHook 中实现了该 Storm 客户端钩子接口。

限制

以下说明适用于该集成的第一个版本。

  • 仅新提交的拓扑会注册到 Atlas,任何生命周期变更都不会反映在 Atlas 中。
  • 提交 Storm 拓扑时,Atlas 服务器必须处于在线状态,元数据才能被采集。
  • 该钩子目前不支持为自定义的 spout 和 bolt 采集血缘信息。

安装

Storm Atlas 钩子需要在客户端手动安装到 Storm 中。

  • 解压 apache-atlas-${project.version}-storm-hook.tar.gz
  • cd apache-atlas-storm-hook-${project.version}
  • 将文件夹 apache-atlas-storm-hook-${project.version}/hook/storm 中的全部内容复制到 $ATLAS_PACKAGE/hook/storm

需要将 $ATLAS_PACKAGE/hook/storm 中的 Storm Atlas hook jar 包复制到 $STORM_HOME/extlib 目录。将 STORM_HOME 替换为 Storm 的安装路径。

安装 Atlas hook 到 Storm 后,请重启所有守护进程。

配置

Storm 配置

Storm Atlas Hook 需要在 Storm 客户端配置文件 $STORM_HOME/conf/storm.yaml 中进行如下配置:

storm.topology.submission.notifier.plugin.class: "org.apache.atlas.storm.hook.StormAtlasHook"

同时设置一个「集群名称」,它将用作 Atlas 中注册对象的命名空间。该名称将用于为 Storm 拓扑、spout 和 bolt 提供命名空间。

其他对象(如数据集)理想情况下应使用生成它们的组件的集群名称来标识。例如,Hive 表和数据库应使用 Hive 中设置的集群名称来标识。如果在客户端提交的 Storm 拓扑 jar 包中提供了 Hive 配置,并且其中定义了集群名称,Storm Atlas 钩子将会获取到该配置。HBase 数据集的情况也类似。如果该配置不可用,则将使用 Storm 配置中设置的集群名称。

atlas.cluster.name: "cluster_name"

在 $STORM_HOME/conf/storm_env.ini 中,按如下方式设置环境变量:

STORM_JAR_JVM_OPTS:"-Datlas.conf=$ATLAS_HOME/conf/"

其中 ATLAS_HOME 指向 Apache Atlas 的安装目录。

你也可以在 Storm 配置中通过编程方式完成此设置:

Config stormConf = new Config();
        ...
        stormConf.put(Config.STORM_TOPOLOGY_SUBMISSION_NOTIFIER_PLUGIN,
                org.apache.atlas.storm.hook.StormAtlasHook.class.getName());

评论

登录后参与评论

正在加载评论…