gpt4 book ai didi

org.apache.gobblin.writer.WatermarkAwareWriter.getUnacknowledgedWatermark()方法的使用及代码示例

转载 作者:知者 更新时间:2024-03-26 10:37:05 29 4
gpt4 key购买 nike

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

WatermarkAwareWriter.getUnacknowledgedWatermark介绍

暂无

代码示例

代码示例来源:origin: apache/incubator-gobblin

public Map<String, CheckpointableWatermark> getUnacknowledgedWatermark() {
 return watermarkAwareWriter.get().getUnacknowledgedWatermark();
}

代码示例来源:origin: apache/incubator-gobblin

@Override
public Map<String, CheckpointableWatermark> getUnacknowledgedWatermark() {
 Preconditions.checkState(isWatermarkCapable());
 return watermarkAwareWriter.get().getUnacknowledgedWatermark();
}

代码示例来源:origin: apache/incubator-gobblin

@Override
public Map<String, CheckpointableWatermark> getUnacknowledgedWatermark() {
 WatermarkTracker watermarkTracker = new MultiWriterWatermarkTracker();
 for (Map.Entry<GenericRecord, DataWriter<D>> entry : this.partitionWriters.asMap().entrySet()) {
  Map<String, CheckpointableWatermark> unacknowledgedWatermark =
    ((WatermarkAwareWriter) entry.getValue()).getUnacknowledgedWatermark();
  if (!unacknowledgedWatermark.isEmpty()) {
   watermarkTracker.unacknowledgedWatermarks(unacknowledgedWatermark);
  }
 }
 return watermarkTracker.getAllUnacknowledgedWatermarks();
}

代码示例来源:origin: apache/incubator-gobblin

@Override
public Map<String, CheckpointableWatermark> getCommittableWatermark() {
 // The committable watermark from a collection of commitable and unacknowledged watermarks is the highest
 // committable watermark that is less than the lowest unacknowledged watermark
 WatermarkTracker watermarkTracker = new MultiWriterWatermarkTracker();
 for (Map.Entry<GenericRecord, DataWriter<D>> entry : this.partitionWriters.asMap().entrySet()) {
  if (entry.getValue() instanceof WatermarkAwareWriter) {
   Map<String, CheckpointableWatermark> commitableWatermarks =
     ((WatermarkAwareWriter) entry.getValue()).getCommittableWatermark();
   if (!commitableWatermarks.isEmpty()) {
    watermarkTracker.committedWatermarks(commitableWatermarks);
   }
   Map<String, CheckpointableWatermark> unacknowledgedWatermark =
     ((WatermarkAwareWriter) entry.getValue()).getUnacknowledgedWatermark();
   if (!unacknowledgedWatermark.isEmpty()) {
    watermarkTracker.unacknowledgedWatermarks(unacknowledgedWatermark);
   }
  }
 }
 return watermarkTracker.getAllCommitableWatermarks(); //TODO: Change this to use List of committables instead
}

代码示例来源:origin: apache/incubator-gobblin

public void testWatermarkComputation(Long committed, Long unacknowledged, Long expected) throws IOException {
 State state = new State();
 state.setProp(ConfigurationKeys.WRITER_PARTITIONER_CLASS, TestPartitioner.class.getCanonicalName());
 String defaultSource = "default";
 WatermarkAwareWriter mockDataWriter = mock(WatermarkAwareWriter.class);
 when(mockDataWriter.isWatermarkCapable()).thenReturn(true);
 when(mockDataWriter.getCommittableWatermark()).thenReturn(Collections.singletonMap(defaultSource,
   new DefaultCheckpointableWatermark(defaultSource, new LongWatermark(committed))));
 when(mockDataWriter.getUnacknowledgedWatermark()).thenReturn(Collections.singletonMap(defaultSource,
   new DefaultCheckpointableWatermark(defaultSource, new LongWatermark(unacknowledged))));
 PartitionAwareDataWriterBuilder builder = mock(PartitionAwareDataWriterBuilder.class);
 when(builder.validatePartitionSchema(any(Schema.class))).thenReturn(true);
 when(builder.forPartition(any(GenericRecord.class))).thenReturn(builder);
 when(builder.withWriterId(any(String.class))).thenReturn(builder);
 when(builder.build()).thenReturn(mockDataWriter);
 PartitionedDataWriter writer = new PartitionedDataWriter<String, String>(builder, state);
 RecordEnvelope<String> recordEnvelope = new RecordEnvelope<String>("0");
 recordEnvelope.addCallBack(
   new AcknowledgableWatermark(new DefaultCheckpointableWatermark(defaultSource, new LongWatermark(0))));
 writer.writeEnvelope(recordEnvelope);
 Map<String, CheckpointableWatermark> watermark = writer.getCommittableWatermark();
 System.out.println(watermark.toString());
 if (expected == null) {
  Assert.assertTrue(watermark.isEmpty(), "Expected watermark to be absent");
 } else {
  Assert.assertTrue(watermark.size() == 1);
  Assert.assertEquals((long) expected, ((LongWatermark) watermark.values().iterator().next().getWatermark()).getValue());
 }
}

代码示例来源:origin: org.apache.gobblin/gobblin-core-base

public Map<String, CheckpointableWatermark> getUnacknowledgedWatermark() {
 return watermarkAwareWriter.get().getUnacknowledgedWatermark();
}

代码示例来源:origin: org.apache.gobblin/gobblin-core-base

@Override
public Map<String, CheckpointableWatermark> getUnacknowledgedWatermark() {
 Preconditions.checkState(isWatermarkCapable());
 return watermarkAwareWriter.get().getUnacknowledgedWatermark();
}

代码示例来源:origin: org.apache.gobblin/gobblin-core

@Override
public Map<String, CheckpointableWatermark> getUnacknowledgedWatermark() {
 WatermarkTracker watermarkTracker = new MultiWriterWatermarkTracker();
 for (Map.Entry<GenericRecord, DataWriter<D>> entry : this.partitionWriters.asMap().entrySet()) {
  Map<String, CheckpointableWatermark> unacknowledgedWatermark =
    ((WatermarkAwareWriter) entry.getValue()).getUnacknowledgedWatermark();
  if (!unacknowledgedWatermark.isEmpty()) {
   watermarkTracker.unacknowledgedWatermarks(unacknowledgedWatermark);
  }
 }
 return watermarkTracker.getAllUnacknowledgedWatermarks();
}

代码示例来源:origin: org.apache.gobblin/gobblin-core

@Override
public Map<String, CheckpointableWatermark> getCommittableWatermark() {
 // The committable watermark from a collection of commitable and unacknowledged watermarks is the highest
 // committable watermark that is less than the lowest unacknowledged watermark
 WatermarkTracker watermarkTracker = new MultiWriterWatermarkTracker();
 for (Map.Entry<GenericRecord, DataWriter<D>> entry : this.partitionWriters.asMap().entrySet()) {
  if (entry.getValue() instanceof WatermarkAwareWriter) {
   Map<String, CheckpointableWatermark> commitableWatermarks =
     ((WatermarkAwareWriter) entry.getValue()).getCommittableWatermark();
   if (!commitableWatermarks.isEmpty()) {
    watermarkTracker.committedWatermarks(commitableWatermarks);
   }
   Map<String, CheckpointableWatermark> unacknowledgedWatermark =
     ((WatermarkAwareWriter) entry.getValue()).getUnacknowledgedWatermark();
   if (!unacknowledgedWatermark.isEmpty()) {
    watermarkTracker.unacknowledgedWatermarks(unacknowledgedWatermark);
   }
  }
 }
 return watermarkTracker.getAllCommitableWatermarks(); //TODO: Change this to use List of committables instead
}

29 4 0
Copyright 2021 - 2024 cfsdn All Rights Reserved 蜀ICP备2022000587号
广告合作:1813099741@qq.com 6ren.com