【发布时间】:2014-07-07 16:55:06
【问题描述】:
我希望能够从分布式(非本地)Storm 拓扑将新条目写入 HBase。有一些 GitHub 项目提供 HBase Mappers 或 pre-made Storm bolts 将元组写入 HBase。这些项目提供了在 LocalCluster 上执行样本的说明。
我在这两个项目中遇到的问题是,直接从 bolt 访问 HBase API,它们都需要将 HBase-site.xml 文件包含在类路径中。使用直接 API 方法,也许还有 GitHub 方法,当您执行 HBaseConfiguration.create(); 时,它会尝试从类路径上的条目中找到它需要的信息。
如何修改风暴螺栓的类路径以包含 Hbase 配置文件?
更新:使用 danehammer 的回答,这就是我的工作方式
将以下文件复制到您的 ~/.storm 目录中:
- hbase-common-0.98.0.2.1.2.0-402-hadoop2.jar
- hbase-site.xml
- storm.yaml :注意:如果您不将storm.yaml 复制到该目录中,则storm jar 命令将不会在类路径中使用该目录(请参阅storm.py python script 以自己查看该逻辑-会很好如果记录在案)
接下来,在拓扑类的 main 方法中获取 HBase 配置并对其进行序列化:
final Configuration hbaseConfig = HBaseConfiguration.create();
final DataOutputBuffer databufHbaseConfig = new DataOutputBuffer();
hbaseConfig.write(databufHbaseConfig);
final byte[] baHbaseConfigSerialized = databufHbaseConfig.getData();
通过构造函数将字节数组传递给您的 spout 类。 spout 类将这个字节数组保存到一个字段中(不要在构造函数中反序列化。我发现如果 spout 有一个配置字段,那么在运行拓扑时会出现无法序列化的异常)
在spout的open方法中,反序列化配置并访问hbase表:
Configuration hBaseConfiguration = new Configuration();
ByteArrayInputStream bas = new ByteArrayInputStream(baHbaseConfigSerialized);
hBaseConfiguration.readFields(new DataInputStream(bas));
HTable tbl = new HTable(hBaseConfiguration, HBASE_TABLE_NAME);
Scan scan = new Scan();
scan.addColumn(Bytes.toBytes("cf"), Bytes.toBytes("YOUR_COLUMN"));
scnrTbl = tbl.getScanner(scan);
现在,在 nextTuple 方法中,您可以使用 Scanner 获取下一行:
Result rsltWaveform = scnrWaveformTbl.next();
从结果中提取你想要的,并将这些值传递给一些可序列化的对象。
【问题讨论】:
-
加一个用于不反序列化构造函数中的字节数组。
标签: java hbase apache-storm