Sam,
I've figured out a reason.
In order to serialize and deserialize a lambda expression the expression
must be referred by a Serializable reference or must be cast to
Serializable. This is a compiler requirement.
To make your code workable you can extend Predicate with your own
interface that extends Serializable as well:
class LambdaStrategyextends Strategy {
private final SerializablePredicate<Integer>predicate;
public LambdaStrategy(SerializablePredicate<Integer> predicate) {
this.predicate = predicate;
}
}
interface SerializablePredicate<T>extends Predicate<T>,*Serializable* {
}
This is not the only one existed solution. Take a look at
ComputeClosureExample source code that is a part of Apache Ignite.
Regards,
Denis
On 9/9/2015 10:33 AM, Sam Adams wrote:
Thanks!
On 9 Sep 2015 07:42, "Denis Magda" <[email protected]
<mailto:[email protected]>> wrote:
Sam,
I'll check this on my side and will return with an answer soon.
Thanks,
Denis
On 9/8/2015 4:51 PM, Sam Adams wrote:
Here's a slightly simplified, self-contained class.
Returning Strategy works, returning LambdaStrategy fails.
Thanks,
Sam
---
package simulator;
import static java.util.stream.Collectors.toList;
import java.util.Collection;
import java.util.List;
import java.util.function.Predicate;
import java.util.stream.IntStream;
import org.apache.ignite.Ignite;
import org.apache.ignite.Ignition;
import org.apache.ignite.cluster.ClusterGroup;
import org.apache.ignite.lang.IgniteCallable;
public class IgniteLambdaTest {
public static void main(String args[]) throws Exception {
new IgniteLambdaTest().run();
}
protected void run() {
try (Ignite ignite =
Ignition.start("C:\\Development\\ignite\\libs\\vrs\\ignite-config.xml"))
{
ClusterGroup remotes = ignite.cluster().forRemotes();
List<TestJob> jobs = IntStream.range(0, 100).mapToObj((x) -> new
TestJob()).collect(toList());
Collection<Strategy> result = ignite.compute(remotes).call(jobs);
result.stream().forEach(System.out::println);
}
}
static class TestJob implements IgniteCallable<Strategy> {
private static final long serialVersionUID = 1L;
@Override
public Strategy call() throws Exception {
//return new Strategy();
return new LambdaStrategy((x) -> false);
}
}
}
class Strategy {
}
class LambdaStrategy extends Strategy {
private final Predicate<Integer> predicate;
public LambdaStrategy(Predicate<Integer> predicate) {
this.predicate = predicate;
}
}
On 8 September 2015 at 14:20, Denis Magda <[email protected]
<mailto:[email protected]>> wrote:
Please, see below
On 9/8/2015 4:03 PM, Sam Adams wrote:
Hi,
I have peer class loading enabled in my XML config:
<property name="peerClassLoadingEnabled" value="true"/>
I assume this is equivalent?
Yes, it is.
Making TestJob static doesn't fix the issue.
Is SimulationJob class also marked with 'static'?
Is LambdaStrategy defined as an inner class? If so, then
please make sure that it's labeled with 'static' as well.
Regards,
Denis
Thanks,
Sam
On 8 September 2015 at 12:37, Denis Magda
<[email protected] <mailto:[email protected]>> wrote:
Hi Sam,
Lambda should work fine.
Please try to do the following and let me know if
anything helps.
- enable peer-class loading feature with
IgniteConfiguration.setPeerClassLoading;
- mark TestJob class with static keyword. This will
avoid serialization of an enclosing class TestJob
belongs to.
Regards,
Denis
On 9/8/2015 2:20 PM, Sam Adams wrote:
Hi,
It looks like I'm unable to return a lambda from a
compute job.
Here is the stack trace:
Exception in thread "main" class
org.apache.ignite.IgniteException: Failed to find class
with given class loader for unmarshalling (make sure
same version of all classes are available on all nodes
or enable peer-class-loading):
sun.misc.Launcher$AppClassLoader@c387f44 at
org.apache.ignite.internal.util.IgniteUtils.convertException(IgniteUtils.java:684)
at
org.apache.ignite.internal.IgniteComputeImpl.call(IgniteComputeImpl.java:352)
at
simulator.IgniteLambdaTest.run(IgniteLambdaTest.java:40) at
simulator.IgniteLambdaTest.main(IgniteLambdaTest.java:31)
Caused by: class
org.apache.ignite.IgniteCheckedException: Failed to
find class with given class loader for unmarshalling
(make sure same version of all classes are available on
all nodes or enable peer-class-loading):
sun.misc.Launcher$AppClassLoader@c387f44 at
org.apache.ignite.marshaller.optimized.OptimizedMarshaller.unmarshal(OptimizedMarshaller.java:242)
at
org.apache.ignite.internal.processors.task.GridTaskWorker.onResponse(GridTaskWorker.java:693)
at
org.apache.ignite.internal.processors.task.GridTaskProcessor.processJobExecuteResponse(GridTaskProcessor.java:932)
at
org.apache.ignite.internal.processors.task.GridTaskProcessor$JobMessageListener.onMessage(GridTaskProcessor.java:1156)
at
org.apache.ignite.internal.managers.communication.GridIoManager.processRegularMessage0(GridIoManager.java:761)
at
org.apache.ignite.internal.managers.communication.GridIoManager.access$1500(GridIoManager.java:59)
at
org.apache.ignite.internal.managers.communication.GridIoManager$5.run(GridIoManager.java:724)
at
java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142)
at
java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617)
at java.lang.Thread.run(Thread.java:745) Caused by:
java.lang.ClassNotFoundException:
simulator.IgniteLambdaTest$SimulationJob$$Lambda$144/739267424
at java.lang.Class.forName0(Native Method) at
java.lang.Class.forName(Class.java:348) at
org.apache.ignite.internal.util.IgniteUtils.forName(IgniteUtils.java:7911)
at
org.apache.ignite.internal.MarshallerContextAdapter.getClass(MarshallerContextAdapter.java:177)
at
org.apache.ignite.marshaller.optimized.OptimizedMarshallerUtils.classDescriptor(OptimizedMarshallerUtils.java:252)
at
org.apache.ignite.marshaller.optimized.OptimizedObjectInputStream.readObjectOverride(OptimizedObjectInputStream.java:246)
at
java.io.ObjectInputStream.readObject(ObjectInputStream.java:365)
at
org.apache.ignite.marshaller.optimized.OptimizedObjectInputStream.readFields(OptimizedObjectInputStream.java:419)
at
org.apache.ignite.marshaller.optimized.OptimizedObjectInputStream.readSerializable(OptimizedObjectInputStream.java:507)
at
org.apache.ignite.marshaller.optimized.OptimizedClassDescriptor.read(OptimizedClassDescriptor.java:751)
at
org.apache.ignite.marshaller.optimized.OptimizedObjectInputStream.readObjectOverride(OptimizedObjectInputStream.java:252)
at
java.io.ObjectInputStream.readObject(ObjectInputStream.java:365)
at
org.apache.ignite.marshaller.optimized.OptimizedMarshaller.unmarshal(OptimizedMarshaller.java:236)
... 9 more
Here is a minimal test case:
protected void run() throws Exception {
try (Ignite ignite = Ignition.start(igniteConfigPath)) {
ClusterGroup remotes = ignite.cluster().forRemotes();
List<SimulationJob> jobs = IntStream.range(0,
100).mapToObj((x) -> new TestJob()).collect(toList());
Collection<LambdaStrategy<?>> result =
ignite.compute(remotes).call(jobs);
result.stream().forEach(System.out::println);
}
}
class TestJob implements
IgniteCallable<LambdaStrategy<?>> {
@Override
public LambdaStrategy<Object> call() throws Exception {
return new LambdaStrategy<Object>((x) -> false);
}
}
Is this a known issue/limitation?
Thanks,
Sam