Spark 2.1写入Ignite 2.2缓存时出现数据节点查找失败错误求助
解决Spark写入Ignite时“Failed to find data nodes for cache”错误
我来帮你捋捋这个问题——这个错误本质是说Ignite集群里找不到能承载你要创建的缓存的数据节点,或者Spark和Ignite的连接、配置环节出了问题。下面是几个实用的排查和解决方向:
1. 先确认Ignite集群的节点角色状态
- 首先得确保你的Ignite集群里至少有一个数据节点(不是客户端节点)在正常运行。你可以用Ignite自带的控制脚本快速检查:
- Linux/macOS:
./control.sh topology - Windows:
control.bat topology
执行后看输出的节点列表,找带有DATA角色的节点。如果全是CLIENT角色,那肯定没法创建缓存——得启动至少一个带数据存储配置的服务端节点才行。
- Linux/macOS:
2. 检查缓存配置的合理性
- 你创建缓存时的配置(比如备份数
backups)要和集群数据节点数量匹配。比如如果设置了backups=2,但集群里只有1个数据节点,Ignite找不到足够的节点存备份,就会报这个错。 - 另外,你的
Custom_Class要确保能被Ignite序列化:要么实现Serializable接口,要么配置Ignite的BinarySerializer来处理这个类的序列化逻辑。
3. 验证Spark与Ignite的连接配置
- 确认
IgniteContext的配置里正确指定了Ignite集群的节点地址,比如:val igniteCfg = new IgniteConfiguration() igniteCfg.setAddresses("192.168.1.100:47500", "192.168.1.101:47500") // 替换成你的Ignite节点实际地址 - 检查防火墙:Spark所在机器和Ignite数据节点之间的默认端口(47500-47509、47100-47109)要开放,不能被防火墙拦截导致通信失败。
4. 调整缓存创建的方式
- 如果你是在Spark里动态创建缓存,可以试试先在Ignite集群的配置文件(
ignite-config.xml)里预定义好这个缓存,再让Spark应用去引用,这样能避免动态创建时的配置不一致问题。比如XML里的缓存配置:
然后在Spark里直接用<bean class="org.apache.ignite.configuration.CacheConfiguration"> <property name="name" value="customCache"/> <property name="indexedTypes"> <list> <value>java.lang.Long</value> <value>com.yourpackage.Custom_Class</value> </list> </property> <property name="backups" value="1"/> </bean>igniteContext.fromCache("customCache")获取缓存即可。
5. 查看Ignite节点日志找细节
- 去Ignite数据节点的
work/log目录下看ignite.log,里面会有更详细的错误堆栈——比如是不是缓存初始化时出了序列化问题,或者节点之间的通信有故障,这些细节能帮你精准定位问题。
这里给你一个修正后的Spark代码示例,供参考:
object Spark_Streaming_Processing { case class Custom_Class( @(QuerySqlField @field)(index = true) a: String, @(QuerySqlField @field)(index = true) b: String, @(QuerySqlField @field)(index = true) c: String, @(QuerySqlField @field)(index = true) d: String, @(QuerySqlField @field)(index = true) e: String ) extends Serializable // 确保实现Serializable接口 def main(args: Array[String]): Unit = { val sparkConf = new SparkConf().setAppName("IgniteSparkWriter") val sc = new SparkContext(sparkConf) // 配置Ignite连接 val igniteCfg = new IgniteConfiguration() igniteCfg.setAddresses("ignite-node-1:47500", "ignite-node-2:47500") // 替换为你的节点地址 // 配置缓存 val cacheCfg = new CacheConfiguration[Long, Custom_Class]("customCache") cacheCfg.setIndexedTypes(classOf[Long], classOf[Custom_Class]) cacheCfg.setBackups(1) // 备份数要小于数据节点数量 val igniteContext = new IgniteContext(sc, () => igniteCfg, true) val cache = igniteContext.getOrCreateCache(cacheCfg) // 后续的写入逻辑... sc.stop() } }
内容的提问来源于stack exchange,提问作者manuel mourato
相关产品推荐
相关产品推荐

