【问题标题】:JanusGraph, Spark cluster failing to connect to CassandraJanusGraph,Spark 集群无法连接到 Cassandra
【发布时间】:2023-03-14 13:25:01
【问题描述】:

我正在尝试在创建 JanusGraph 的集群上运行 Spark 作业。

我有一个 JanusGraph 服务器实例,Cassandra,ES 在单台机器上运行,只有 Spark 计算发生在集群上。 (基本上我在机器上做了一个janusgraph.sh start

我的配置如下(x是我运行上述实例的机器的IP):

def getGraph(): JanusGraph = {
    val config = JanusGraphFactory.build()
    config.set("storage.backend", "cassandrathrift")
    config.set("storage.cassandrathrift.keyspace", "jgex")
    config.set("storage.hostname", "x")
    config.set("index.jgex.backend", "elasticsearch")
    config.set("index.jgex.index-name", "jgex")
    config.set("jgex.hostname", "x")
    config.open()
  }

但是当我对集群上的胖罐子做一个spark-submit 时,我得到了这个:

    java.lang.IllegalArgumentException: Could not instantiate implementation: org.janusgraph.diskstorage.cassandra.thrift.CassandraThriftStoreManager
    at org.janusgraph.util.system.ConfigurationUtil.instantiate(ConfigurationUtil.java:69)
    at org.janusgraph.diskstorage.Backend.getImplementationClass(Backend.java:477)
    at org.janusgraph.diskstorage.Backend.getStorageManager(Backend.java:409)
    at org.janusgraph.graphdb.configuration.GraphDatabaseConfiguration.<init>(GraphDatabaseConfiguration.java:1376)
    at org.janusgraph.core.JanusGraphFactory.open(JanusGraphFactory.java:164)
    at org.janusgraph.core.JanusGraphFactory.open(JanusGraphFactory.java:133)
    at org.janusgraph.core.JanusGraphFactory.open(JanusGraphFactory.java:123)
    at org.janusgraph.core.JanusGraphFactory$Builder.open(JanusGraphFactory.java:264)
    at janus_create$.getGraph(janus_create.scala:66)
    at janus_create$.makePropertiesandIndexes(janus_create.scala:830)
    at janus_create$.main(janus_create.scala:921)
    at janus_create.main(janus_create.scala)
    at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
    at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
    at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
    at java.lang.reflect.Method.invoke(Method.java:498)
    at org.apache.spark.deploy.yarn.ApplicationMaster$$anon$2.run(ApplicationMaster.scala:627)
Caused by: java.lang.reflect.InvocationTargetException
    at sun.reflect.NativeConstructorAccessorImpl.newInstance0(Native Method)
    at sun.reflect.NativeConstructorAccessorImpl.newInstance(NativeConstructorAccessorImpl.java:62)
    at sun.reflect.DelegatingConstructorAccessorImpl.newInstance(DelegatingConstructorAccessorImpl.java:45)
    at java.lang.reflect.Constructor.newInstance(Constructor.java:423)
    at org.janusgraph.util.system.ConfigurationUtil.instantiate(ConfigurationUtil.java:58)
    ... 16 more
Caused by: org.janusgraph.diskstorage.TemporaryBackendException: Temporary failure in storage backend
    at org.janusgraph.diskstorage.cassandra.thrift.CassandraThriftStoreManager.getCassandraPartitioner(CassandraThriftStoreManager.java:219)
    at org.janusgraph.diskstorage.cassandra.thrift.CassandraThriftStoreManager.<init>(CassandraThriftStoreManager.java:198)
    ... 21 more
Caused by: org.apache.thrift.transport.TTransportException: java.net.ConnectException: Connection refused (Connection refused)
    at org.apache.thrift.transport.TSocket.open(TSocket.java:187)
    at org.apache.thrift.transport.TFramedTransport.open(TFramedTransport.java:81)
    at org.janusgraph.diskstorage.cassandra.thrift.thriftpool.CTConnectionFactory.makeRawConnection(CTConnectionFactory.java:110)
    at org.janusgraph.diskstorage.cassandra.thrift.thriftpool.CTConnectionFactory.makeObject(CTConnectionFactory.java:74)
    at org.janusgraph.diskstorage.cassandra.thrift.thriftpool.CTConnectionFactory.makeObject(CTConnectionFactory.java:43)
    at org.apache.commons.pool.impl.GenericKeyedObjectPool.borrowObject(GenericKeyedObjectPool.java:1179)
    at org.janusgraph.diskstorage.cassandra.thrift.CassandraThriftStoreManager.getCassandraPartitioner(CassandraThriftStoreManager.java:216)
    ... 22 more
Caused by: java.net.ConnectException: Connection refused (Connection refused)
    at java.net.PlainSocketImpl.socketConnect(Native Method)
    at java.net.AbstractPlainSocketImpl.doConnect(AbstractPlainSocketImpl.java:350)
    at java.net.AbstractPlainSocketImpl.connectToAddress(AbstractPlainSocketImpl.java:206)
    at java.net.AbstractPlainSocketImpl.connect(AbstractPlainSocketImpl.java:188)
    at java.net.SocksSocketImpl.connect(SocksSocketImpl.java:392)
    at java.net.Socket.connect(Socket.java:589)
    at org.apache.thrift.transport.TSocket.open(TSocket.java:182)
    ... 28 more

我尝试在 cassandra 和 cassandrathrift 之间切换,但两者都不起作用。另外,我在哪里指定我的 gremlin 在哪里运行。这有关系吗?

【问题讨论】:

  • 您是为单个 Cassandra 节点使用预打包的 JanusGraph 发行版,还是使用独立的 Cassandra 节点?澄清一下,你有一个 Spark 集群在与 Cassandra/Elasticsearch/JanusGraph 不同的机器上运行?
  • 您好,很抱歉回复晚了。我已经弄清楚发生了什么,cassandra 和 ES 没有响应远程请求。我正在运行预打包的 JanusGraph 发行版,是的,我有一个 Spark 集群在不同的机器上运行。远程设置 ES/Cassandra 需要一些特殊配置吗?另外,如何运行 Janusgraph 服务器的多个实例?谢谢。
  • 我尝试修复了很多东西,我猜与此相关的稀疏文档让事情变得比应有的更困难。 Janusgraph 服务器本质上是一个 gremlin 服务器,我可以进行 conf 更改以使 gremlin 服务器指向正确的 cassandra 和 ES 实例。但是如何配置我的 spark 作业以指向正确的 gremlin/janusgraph 服务器?提前致谢。

标签: apache-spark cassandra gremlin janusgraph


【解决方案1】:

预打包的发行版假定每个 Cassandra、Elasticsearch 和 Gremlin 服务器都有一个 localhost 节点。注意堆栈跟踪底部的java.net.ConnectException: Connection refused。如果您在远程集群上运行 Spark,则需要确保服务器在非本地主机地址上可用。

  • 首先使用bin/janusgraph.sh stop 停止分发
  • 使用机器的IP地址(Cassandra docs)更新$JANUSGRAPH_HOME/conf/cassandra/cassandra.yaml中的listen_addressrpc_address

  • 使用机器的IP地址将network.host添加到$JANUSGRAPH_HOME/elasticsearch/config/elasticsearch.yml (Elasticsearch docs)

  • 使用机器的IP地址(JanusGraph docs)更新$JANUSGRAPH_HOME/conf/gremlin-server/gremlin-server.yaml中的host

  • 假设您使用默认 Gremlin 服务器配置 gremlin-server.yaml,您需要使用机器的 IP 地址更新 $JANUSGRAPH_HOME/conf/gremlin-server/janusgraph-cassandra-es-server.properties 中的属性文件。使用机器的 IP 地址更新 storage.hostnameindex.search.hostname,匹配上面的服务器设置。

更新:您的图形连接属性中似乎存在可能导致问题的错误:

  • config.set("storage.cassandrathrift.keyspace", "jgex") -- 应该是"storage.cassandra.keyspace"
  • config.set("jgex.hostname", "x") -- 应该是 "index.jgex.hostname"

【讨论】:

  • 谢谢,我做了你提到的一切,现在它抛出了一个类似的错误,只是针对 ElasticSearch。另外,我想问一下,据我所知,spark 作业应该与 Gremlin 服务器通信,而它们又应该与后端/索引后端通信。但是在我的 spark 代码中,我需要指定 cassandra 和 ES 的位置。那么,火花工作是如何知道我的小精灵在哪里的呢?我需要在我的 spark 作业中指定它吗?
  • 请查看更新后的答案。您上面的代码直接连接到存储后端和索引后端,这是一种有效的方法。如果您想通过 Gremlin 服务器进行连接,则需要使用 Gremlin Driver 方法。
  • 谢谢,效果很好!另外,如果我想添加一个 Array[Double] 属性,如何在 Scala 中定义它? mgmt.makePropertyKey("prop1").dataType(classOf[java.util.Arrays]).make() 不起作用。
猜你喜欢
  • 1970-01-01
  • 2018-04-10
  • 2021-11-14
  • 2017-10-03
  • 2014-06-14
  • 2017-08-19
  • 2017-01-30
  • 2016-04-25
  • 1970-01-01
相关资源
最近更新 更多