Storm
Apache Storm 的 Apache Atlas 钩子
简介
Apache Storm 是一个分布式实时计算系统。Storm 使得可靠地处理无界数据流变得容易,它之于实时处理,就如同 Hadoop 之于批处理。其处理过程本质上是由节点构成的 DAG,称为 拓扑(topology)。
Apache Atlas 是一个元数据仓库,能够实现端到端的数据血缘、搜索以及业务分类的关联。
该集成的目标是将运行时的拓扑元数据连同底层的数据源、目标、派生过程以及任何可用的业务上下文一并推送,以便 Atlas 能够捕获该拓扑的血缘关系。
此过程包含以下两个部分,详见下文:
- 用于表示 Storm 中各概念的数据模型
- 用于在 Atlas 中更新元数据的 Storm Atlas 钩子
Storm 数据模型
数据模型在 Atlas 中以类型(Types)的形式表示。它包含拓扑图中各类节点的描述,例如 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 拓扑的客户端上提供了钩子 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 中通过编程方式这样设置:
Config stormConf = new Config();
...
stormConf.put(Config.STORM_TOPOLOGY_SUBMISSION_NOTIFIER_PLUGIN,
org.apache.atlas.storm.hook.StormAtlasHook.class.getName());评论
登录后参与评论
KnowForge