本文整理了Java中com.hazelcast.jet.impl.util.Util.uncheckCall()
方法的一些代码示例,展示了Util.uncheckCall()
的具体用法。这些代码示例主要来源于Github
/Stackoverflow
/Maven
等平台,是从一些精选项目中提取出来的代码,具有较强的参考意义,能在一定程度帮忙到你。Util.uncheckCall()
方法的具体详情如下:
包路径:com.hazelcast.jet.impl.util.Util
类名称:Util
方法名:uncheckCall
暂无
代码示例来源:origin: hazelcast/hazelcast-jet
@Override
public boolean complete() {
return uncheckCall(this::tryComplete);
}
代码示例来源:origin: hazelcast/hazelcast-jet
@Override
public AttributeList getAttributes(String[] attributes) {
return Arrays.stream(attributes)
.filter(metrics::containsKey)
.map(a -> uncheckCall(() -> new Attribute(a, getAttribute(a))))
.collect(toCollection(AttributeList::new));
}
代码示例来源:origin: hazelcast/hazelcast-jet
@Override
public boolean complete() {
if (traverser == null) {
resultSet = uncheckCall(() -> resultSetFn.createResultSet(connection, parallelism, index));
traverser = ((Traverser<ResultSet>) () -> uncheckCall(() -> resultSet.next() ? resultSet : null))
.map(mapOutputFn);
}
return emitFromTraverser(traverser);
}
代码示例来源:origin: hazelcast/hazelcast-jet
@Override
@Nonnull
public List<Processor> get(int count) {
Map<Integer, List<IndexedInputSplit>> processorToSplits =
range(0, assignedSplits.size()).mapToObj(i -> new SimpleImmutableEntry<>(i, assignedSplits.get(i)))
.collect(groupingBy(e -> e.getKey() % count,
mapping(Entry::getValue, toList())));
range(0, count)
.forEach(processor -> processorToSplits.computeIfAbsent(processor, x -> emptyList()));
InputFormat inputFormat = jobConf.getInputFormat();
return processorToSplits
.values().stream()
.map(splits -> splits.isEmpty()
? Processors.noopP().get()
: new ReadHdfsP<>(splits.stream()
.map(IndexedInputSplit::getSplit)
.map(split -> uncheckCall(() ->
inputFormat.getRecordReader(split, jobConf, NULL)))
.collect(toList()), mapper)
).collect(toList());
}
代码示例来源:origin: com.hazelcast.jet/hazelcast-jet-hadoop
@Override
@Nonnull
public List<Processor> get(int count) {
Map<Integer, List<IndexedInputSplit>> processorToSplits =
range(0, assignedSplits.size()).mapToObj(i -> new SimpleImmutableEntry<>(i, assignedSplits.get(i)))
.collect(groupingBy(e -> e.getKey() % count,
mapping(Entry::getValue, toList())));
range(0, count)
.forEach(processor -> processorToSplits.computeIfAbsent(processor, x -> emptyList()));
InputFormat inputFormat = jobConf.getInputFormat();
return processorToSplits
.values().stream()
.map(splits -> splits.isEmpty()
? Processors.noopP().get()
: new ReadHdfsP<>(splits.stream()
.map(IndexedInputSplit::getSplit)
.map(split -> uncheckCall(() ->
inputFormat.getRecordReader(split, jobConf, NULL)))
.collect(toList()), mapper)
).collect(toList());
}
代码示例来源:origin: hazelcast/hazelcast-jet
@Override
protected Class<?> findClass(String name) throws ClassNotFoundException {
if (isEmpty(name)) {
return null;
}
InputStream classBytesStream = resourceStream(name.replace('.', '/') + ".class");
if (classBytesStream == null) {
throw new ClassNotFoundException(name + ". Add it using " + JobConfig.class.getSimpleName()
+ " or start all members with it on classpath");
}
byte[] classBytes = uncheckCall(() -> IOUtil.toByteArray(classBytesStream));
return defineClass(name, classBytes, 0, classBytes.length);
}
代码示例来源:origin: hazelcast/hazelcast-jet
private boolean isSplitLocalForMember(InputSplit split, Address memberAddr) {
try {
final InetAddress inetAddr = memberAddr.getInetAddress();
return Arrays.stream(split.getLocations())
.flatMap(loc -> Arrays.stream(uncheckCall(() -> InetAddress.getAllByName(loc))))
.anyMatch(inetAddr::equals);
} catch (IOException e) {
if (e instanceof UnknownHostException) {
logger.warning("Failed to resolve host name for the split, " +
"will use host name equality to determine data locality", e);
return isSplitLocalForMember(split, memberAddr.getScopedHost());
}
throw sneakyThrow(e);
}
}
代码示例来源:origin: com.hazelcast.jet/hazelcast-jet-hadoop
private boolean isSplitLocalForMember(InputSplit split, Address memberAddr) {
try {
final InetAddress inetAddr = memberAddr.getInetAddress();
return Arrays.stream(split.getLocations())
.flatMap(loc -> Arrays.stream(uncheckCall(() -> InetAddress.getAllByName(loc))))
.anyMatch(inetAddr::equals);
} catch (IOException e) {
if (e instanceof UnknownHostException) {
logger.warning("Failed to resolve host name for the split, " +
"will use host name equality to determine data locality", e);
return isSplitLocalForMember(split, memberAddr.getScopedHost());
}
throw sneakyThrow(e);
}
}
代码示例来源:origin: hazelcast/hazelcast-jet
@Override
protected void init(@Nonnull Context context) {
session = sessionFn.apply(connection);
consumer = consumerFn.apply(session);
traverser = ((Traverser<Message>) () -> uncheckCall(() -> consumer.receiveNoWait()))
.flatMap(t -> eventTimeMapper.flatMapEvent(projectionFn.apply(t), 0, handleJmsTimestamp(t)))
.peek(item -> flushFn.accept(session));
}
代码示例来源:origin: hazelcast/hazelcast-jet
private Traverser<Object> traverser(byte[] data) {
BufferObjectDataInput in = serializationService.createObjectDataInput(data);
return () -> uncheckCall(() -> {
Object key = in.readObject();
if (key == SnapshotDataValueTerminator.INSTANCE) {
in.close();
return null;
}
Object value = in.readObject();
return key instanceof BroadcastKey
? new BroadcastEntry(key, value)
: entry(key, value);
});
}
内容来源于网络,如有侵权,请联系作者删除!