com.hazelcast.jet.impl.util.Util.uncheckCall()方法的使用及代码示例

x33g5p2x  于2022-02-01 转载在 其他  
字(5.1k)|赞(0)|评价(0)|浏览(235)

本文整理了Java中com.hazelcast.jet.impl.util.Util.uncheckCall()方法的一些代码示例,展示了Util.uncheckCall()的具体用法。这些代码示例主要来源于Github/Stackoverflow/Maven等平台,是从一些精选项目中提取出来的代码,具有较强的参考意义,能在一定程度帮忙到你。Util.uncheckCall()方法的具体详情如下:
包路径:com.hazelcast.jet.impl.util.Util
类名称:Util
方法名:uncheckCall

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);
  });
}

相关文章

微信公众号

最新文章

更多