【发布时间】:2014-03-12 12:19:48
【问题描述】:
在本地集群中运行拓扑后,我创建了一个远程风暴集群(storm-deploy Nathan)。在创建具有“包依赖项”的可运行 jar 之前,我已经从 eclipse 的构建路径中删除了 Storm jar。我的拓扑使用storm-kafka-0.9.0-wip16a-scala292.jar,我要么留在构建路径中并在创建可运行jar之前从构建路径中删除(只是为了尝试解决这个问题..)。当我使用以下命令时:
./storm jar /home/ubuntu/Virtual/stormTopologia4.jar org.vicomtech.main.StormTopologia
它总是回复:
Exception in thread "main" java.lang.NoClassDefFoundError: OpaqueTridentKafkaSpout
at java.lang.Class.getDeclaredMethods0(Native Method)
at java.lang.Class.privateGetDeclaredMethods(Class.java:2451)
at java.lang.Class.getMethod0(Class.java:2694)
at java.lang.Class.getMethod(Class.java:1622)
at sun.launcher.LauncherHelper.getMainMethod(LauncherHelper.java:494)
at sun.launcher.LauncherHelper.checkAndLoadMain(LauncherHelper.java:486)
Caused by: java.lang.ClassNotFoundException: OpaqueTridentKafkaSpout
at java.net.URLClassLoader$1.run(URLClassLoader.java:366)
at java.net.URLClassLoader$1.run(URLClassLoader.java:355)
at java.security.AccessController.doPrivileged(Native Method)
at java.net.URLClassLoader.findClass(URLClassLoader.java:354)
at java.lang.ClassLoader.loadClass(ClassLoader.java:423)
at sun.misc.Launcher$AppClassLoader.loadClass(Launcher.java:308)
at java.lang.ClassLoader.loadClass(ClassLoader.java:356)
由于此拓扑在 AWS 上作为可运行 jar 的单个实例上运行良好,我无法弄清楚我缺少什么... 这是我的主要方法中的代码:
Config conf = new Config();
OpaqueTridentKafkaSpout tridentSpout = crearSpout(
kafkadir, "test");
OpaqueTridentKafkaSpout logUpvSpout = crearSpout(kafkadir,
"logsUpv");
OpaqueTridentKafkaSpout logSnortSpout = crearSpout(
kafkadir, "logsSnort");
try {
StormSubmitter.submitTopology(
"hackaton",
conf,
buildTopology( tridentSpout, logUpvSpout,
logSnortSpout));
} catch (AlreadyAliveException | InvalidTopologyException e) {
e.printStackTrace();
}
} catch (IOException e) {
e.printStackTrace();
} catch (TwitterException e) {
e.printStackTrace();
}
}
private static OpaqueTridentKafkaSpout crearSpout(
String testKafkaBrokerHost, String topic) {
KafkaConfig.ZkHosts hosts = new ZkHosts(testKafkaBrokerHost, "/brokers");
TridentKafkaConfig config = new TridentKafkaConfig(hosts, topic);
config.forceStartOffsetTime(-2);
config.scheme = new SchemeAsMultiScheme(new StringScheme());
return new OpaqueTridentKafkaSpout(config);
}
public static StormTopology buildTopology(OpaqueTridentKafkaSpout tridentSpout,
OpaqueTridentKafkaSpout logUpvSpout,
OpaqueTridentKafkaSpout logSnortSpout
) throws IOException,
TwitterException {
TridentTopology topology = new TridentTopology();
topology.newStream("tweets2", tridentSpout)
.each(new Fields("str"), new OnlyEnglishSpanish())
.each(new Fields("str"), new WholeTweetToMongo())
.each(new Fields("str"), new TextLangExtracter(),
new Fields("text", "lang")).parallelismHint(6)
.project(new Fields("text", "lang"))
.partitionBy(new Fields("lang"))
.each(new Fields("text", "lang"), new Analisis(),
new Fields("result")).parallelismHint(6)
.each(new Fields("result"), new ResultToMongo());
return topology.build();
}
有什么方法可以让 OpaqueTridentKafkaSpout 可用? 提前谢谢你
希望这不是一个愚蠢的暗示,我是这个领域的新手
【问题讨论】:
标签: java amazon-web-services apache-storm apache-kafka