【问题标题】:Apache Ignite: Serialization error related to Data StreamingApache Ignite:与数据流相关的序列化错误
【发布时间】:2016-09-16 20:45:08
【问题描述】:

我正在尝试研究 Apache Ignite 流式传输的工作原理。我有 2 个节点集群设置(都在本地主机上),并且我启动了一个客户端节点,它使用 StreamTransformer 和 EntryProcessor 运行流代码。结果,在我的一个节点中,我无法反序列化异常。我的代码是 Ignite 文档中简化的 WordCount 示例:

public class StreamingExample {`
public static class StreamingExampleCacheEntryProcessor implements CacheEntryProcessor<String, Long, Object> {
    @Override
    public Object process(MutableEntry<String, Long> e, Object... arg) throws EntryProcessorException {
        Long val = e.getValue();
        e.setValue(val == null ? 1L : val + 1);
        return null;
    }
}

public static void main(String[] args) throws IgniteException, IOException {
    Ignition.setClientMode(true);
    try (Ignite ignite = Ignition.start("examples/config/example-ignite.xml")) {
        IgniteCache<String, Long> stmCache = ignite.getOrCreateCache("mycache");
        try (IgniteDataStreamer<String, Long> stmr = ignite.dataStreamer(stmCache.getName())) {
            stmr.allowOverwrite(true);
            stmr.receiver(StreamTransformer.from(new StreamingExampleCacheEntryProcessor()));
            stmr.addData("word", 1L);
            System.out.println("Finished");
        }
    }
}

}

例外我得到两个节点之一是

[23:38:23] 拓扑快照 [ver=5,servers=2,clients=1,CPUs=4,heap=3.3GB] 线程“pub-#9%null%”类 org.apache.ignite.binary.BinaryObjectException 中的异常:无法使用优化的编组器解组对象 在 org.apache.ignite.internal.binary.BinaryUtils.doReadOptimized(BinaryUtils.java:1595) 在 org.apache.ignite.internal.binary.BinaryReaderExImpl.deserialize(BinaryReaderExImpl.java:1663) 在 org.apache.ignite.internal.binary.GridBinaryMarshaller.deserialize(GridBinaryMarshaller.java:298) 在 org.apache.ignite.internal.binary.BinaryMarshaller.unmarshal(BinaryMarshaller.java:109) 在 org.apache.ignite.internal.processors.datastreamer.DataStreamProcessor.processRequest(DataStreamProcessor.java:278) 在 org.apache.ignite.internal.processors.datastreamer.DataStreamProcessor.access$000(DataStreamProcessor.java:50) 在 org.apache.ignite.internal.processors.datastreamer.DataStreamProcessor$1.onMessage(DataStreamProcessor.java:80) 在 org.apache.ignite.internal.managers.communication.GridIoManager.invokeListener(GridIoManager.java:1238) 在 org.apache.ignite.internal.managers.communication.GridIoManager.processRegularMessage0(GridIoManager.java:866) 在 org.apache.ignite.internal.managers.communication.GridIoManager.access $1700(GridIoManager.java:106) 在 org.apache.ignite.internal.managers.communication.GridIoManager$5.run(GridIoManager.java:829) 在 java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1145) 在 java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:615) 在 java.lang.Thread.run(Thread.java:745) 原因:类 org.apache.ignite.IgniteCheckedException:无法找到具有给定类加载器的类以进行解组(确保所有类的相同版本可用 在所有节点上标记或启用对等类加载):sun.misc.Launcher$AppClassLoader@4e857327 在 org.apache.ignite.marshaller.optimized.OptimizedMarshaller.unmarshal(OptimizedMarshaller.java:224) 在 org.apache.ignite.internal.binary.BinaryUtils.doReadOptimized(BinaryUtils.java:1592) ... 13 更多 引起:java.lang.ClassNotFoundException:gridgaingames.StreamingExample$StreamingExampleCacheEntryProcessor 在 java.net.URLClassLoader$1.run(URLClassLoader.java:366) 在 java.net.URLClassLoader$1.run(URLClassLoader.java:355) 在 java.security.AccessController.doPrivileged(本机方法) 在 java.net.URLClassLoader.findClass(URLClassLoader.java:354) 在 java.lang.ClassLoader.loadClass(ClassLoader.java:425) 在 sun.misc.Launcher$AppClassLoader.loadClass(Launcher.java:308) 在 java.lang.ClassLoader.loadClass(ClassLoader.java:358) 在 java.lang.Class.forName0(本机方法) 在 java.lang.Class.forName(Class.java:274) 在 org.apache.ignite.internal.util.IgniteUtils.forName(IgniteUtils.java:8350) 在 org.apache.ignite.internal.MarshallerContextAdapter.getClass(MarshallerContextAdapter.java:185) 在 org.apache.ignite.marshaller.optimized.OptimizedMarshallerUtils.classDescriptor(OptimizedMarshallerUtils.java:266) 在 org.apache.ignite.marshaller.optimized.OptimizedObjectInputStream.readObjectOverride(OptimizedObjectInputStream.java:318) 在 java.io.ObjectInputStream.readObject(ObjectInputStream.java:364) 在 org.apache.ignite.marshaller.optimized.OptimizedObjectInputStream.readFields(OptimizedObjectInputStream.java:491) 在 org.apache.ignite.marshaller.optimized.OptimizedObjectInputStream.readSerializable(OptimizedObjectInputStream.java:579) 在 org.apache.ignite.marshaller.optimized.OptimizedClassDescriptor.read(OptimizedClassDescriptor.java:841) 在 org.apache.ignite.marshaller.optimized.OptimizedObjectInputStream.readObjectOverride(OptimizedObjectInputStream.java:324) 在 java.io.ObjectInputStream.readObject(ObjectInputStream.java:364) 在 org.apache.ignite.marshaller.optimized.OptimizedMarshaller.unmarshal(OptimizedMarshaller.java:218) ... 14 更多

有几样东西我得不到。

1) 我该如何解决?

2) 由于这不是“广播”之类的,我认为 Ignite 仅在调用节点上运行流式处理代码。看来我错了。那么我的 Streaming 代码在哪里执行呢?

3) 打印“已完成”行后,我的代码不会停止。为什么?看起来一些非守护线程仍然存在。这是阻止我的客户端节点退出的流代码吗?

PS

对等类加载已启用。如果我运行一些在许多节点上执行代码的广播示例 - 它可以正常工作。

【问题讨论】:

    标签: java classloader gridgain ignite


    【解决方案1】:

    基本上IgniteDataStreamer 在发送方(在您的示例中为客户端)准备数据批次,并立即将它们发送到应存储特定键值元组的目标节点。牢记这一点,您的问题的答案如下:

    1. 在将条目放入缓存之前,转换器会在目标节点(服务器节点)上执行。这意味着服务器节点的类路径中必须有转换器的类,或者,您必须启用对等类加载。就个人而言,后者是更灵活、更可取的解决方案。
    2. 正如上面所解释的,发送者只需准备发送到部署缓存的所有服务器的批处理。服务器只接收那些包含元组的批次,这些元组服务器是主服务器或备份服务器。
    3. 批次的刷新发生在后台,因为IgniteDataStreamer 用于快速数据预加载或复杂流处理 (CEP)。有许多参数可以让您调整刷新 - autoFlustFrequencyperNodeBufferSize

    最后,对于预加载需求(当缓存为空并且您需要填充它们时)我建议将allowOverwrite 设置为false,这将允许流媒体分别为主节点和备份节点准备和发送批次。如果此参数设置为true,则批次仅在主节点上发送,主节点在更新其数据版本和相应备份之后,通过基本cache.put 操作注入数据。如果您只需要预加载缓存,这种方法会比较慢。

    【讨论】:

    • 谢谢,dmagda!但我找不到关于 ClassNotFoundException 的主要问题的答案。 PeerToPeer 类加载已启用(如您所见,我正在使用 Ignite 附带的配置之一 - 示例/配置/示例-ignite.xml)。我仍然收到此错误。我猜测可能是 Transformer 的代码在 Transformer 的类被节点实际接收之前启动。这是根本原因,还是别的什么?以及如何解决?
    • 我设法重现了这个问题 (issues.apache.org/jira/browse/IGNITE-3935)。感谢举报!请参阅工单中建议的解决方法。
    猜你喜欢
    • 2016-07-02
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-03-02
    • 2018-09-02
    • 2019-01-05
    相关资源
    最近更新 更多