org.apache.flink.streaming.api.operators.Output.emitLatencyMarker()方法的使用及代码示例

x33g5p2x  于2022-01-26 转载在 其他  
字(5.9k)|赞(0)|评价(0)|浏览(101)

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

Output.emitLatencyMarker介绍

暂无

代码示例

代码示例来源:origin: apache/flink

@Override
public void emitLatencyMarker(LatencyMarker latencyMarker) {
  output.emitLatencyMarker(latencyMarker);
}

代码示例来源:origin: apache/flink

@Override
public void emitLatencyMarker(LatencyMarker latencyMarker) {
  if (outputs.length <= 0) {
    // ignore
  } else if (outputs.length == 1) {
    outputs[0].emitLatencyMarker(latencyMarker);
  } else {
    // randomly select an output
    outputs[random.nextInt(outputs.length)].emitLatencyMarker(latencyMarker);
  }
}

代码示例来源:origin: apache/flink

@Override
public void emitLatencyMarker(LatencyMarker latencyMarker) {
  // randomly select an output
  allOutputs[random.nextInt(allOutputs.length)].emitLatencyMarker(latencyMarker);
}

代码示例来源:origin: apache/flink

protected void reportOrForwardLatencyMarker(LatencyMarker marker) {
  // all operators are tracking latencies
  this.latencyStats.reportLatency(marker);
  // everything except sinks forwards latency markers
  this.output.emitLatencyMarker(marker);
}

代码示例来源:origin: apache/flink

@Override
  public void onProcessingTime(long timestamp) throws Exception {
    try {
      // ProcessingTimeService callbacks are executed under the checkpointing lock
      output.emitLatencyMarker(new LatencyMarker(timestamp, operatorId, subtaskIndex));
    } catch (Throwable t) {
      // we catch the Throwables here so that we don't trigger the processing
      // timer services async exception handler
      LOG.warn("Error while emitting latency marker.", t);
    }
  }
},

代码示例来源:origin: com.alibaba.blink/flink-table

@Override
  public void emitLatencyMarker(LatencyMarker latencyMarker) {
    rawOutput.emitLatencyMarker(latencyMarker);
  }
};

代码示例来源:origin: org.apache.flink/flink-streaming-java

@Override
public void emitLatencyMarker(LatencyMarker latencyMarker) {
  output.emitLatencyMarker(latencyMarker);
}

代码示例来源:origin: org.apache.flink/flink-streaming-java_2.11

@Override
public void emitLatencyMarker(LatencyMarker latencyMarker) {
  output.emitLatencyMarker(latencyMarker);
}

代码示例来源:origin: org.apache.flink/flink-streaming-java_2.11

@Override
public void emitLatencyMarker(LatencyMarker latencyMarker) {
  if (outputs.length <= 0) {
    // ignore
  } else if (outputs.length == 1) {
    outputs[0].emitLatencyMarker(latencyMarker);
  } else {
    // randomly select an output
    outputs[random.nextInt(outputs.length)].emitLatencyMarker(latencyMarker);
  }
}

代码示例来源:origin: org.apache.flink/flink-streaming-java_2.10

@Override
public void emitLatencyMarker(LatencyMarker latencyMarker) {
  if (outputs.length <= 0) {
    // ignore
  } else if (outputs.length == 1) {
    outputs[0].emitLatencyMarker(latencyMarker);
  } else {
    // randomly select an output
    outputs[random.nextInt(outputs.length)].emitLatencyMarker(latencyMarker);
  }
}

代码示例来源:origin: org.apache.flink/flink-streaming-java

@Override
public void emitLatencyMarker(LatencyMarker latencyMarker) {
  if (outputs.length <= 0) {
    // ignore
  } else if (outputs.length == 1) {
    outputs[0].emitLatencyMarker(latencyMarker);
  } else {
    // randomly select an output
    outputs[random.nextInt(outputs.length)].emitLatencyMarker(latencyMarker);
  }
}

代码示例来源:origin: org.apache.flink/flink-streaming-java_2.11

@Override
public void emitLatencyMarker(LatencyMarker latencyMarker) {
  // randomly select an output
  allOutputs[random.nextInt(allOutputs.length)].emitLatencyMarker(latencyMarker);
}

代码示例来源:origin: org.apache.flink/flink-streaming-java_2.10

@Override
public void emitLatencyMarker(LatencyMarker latencyMarker) {
  // randomly select an output
  allOutputs[random.nextInt(allOutputs.length)].emitLatencyMarker(latencyMarker);
}

代码示例来源:origin: org.apache.flink/flink-streaming-java

@Override
public void emitLatencyMarker(LatencyMarker latencyMarker) {
  // randomly select an output
  allOutputs[random.nextInt(allOutputs.length)].emitLatencyMarker(latencyMarker);
}

代码示例来源:origin: org.apache.flink/flink-streaming-java_2.11

protected void reportOrForwardLatencyMarker(LatencyMarker marker) {
  // all operators are tracking latencies
  this.latencyStats.reportLatency(marker);
  // everything except sinks forwards latency markers
  this.output.emitLatencyMarker(marker);
}

代码示例来源:origin: org.apache.flink/flink-streaming-java_2.10

@Override
  public void onProcessingTime(long timestamp) throws Exception {
    try {
      // ProcessingTimeService callbacks are executed under the checkpointing lock
      output.emitLatencyMarker(new LatencyMarker(timestamp, vertexID, subtaskIndex));
    } catch (Throwable t) {
      // we catch the Throwables here so that we don't trigger the processing
      // timer services async exception handler
      LOG.warn("Error while emitting latency marker.", t);
    }
  }
},

代码示例来源:origin: org.apache.flink/flink-streaming-java

@Override
  public void onProcessingTime(long timestamp) throws Exception {
    try {
      // ProcessingTimeService callbacks are executed under the checkpointing lock
      output.emitLatencyMarker(new LatencyMarker(timestamp, operatorId, subtaskIndex));
    } catch (Throwable t) {
      // we catch the Throwables here so that we don't trigger the processing
      // timer services async exception handler
      LOG.warn("Error while emitting latency marker.", t);
    }
  }
},

代码示例来源:origin: org.apache.flink/flink-streaming-java

protected void reportOrForwardLatencyMarker(LatencyMarker marker) {
  // all operators are tracking latencies
  this.latencyStats.reportLatency(marker);
  // everything except sinks forwards latency markers
  this.output.emitLatencyMarker(marker);
}

代码示例来源:origin: org.apache.flink/flink-streaming-java_2.10

protected void reportOrForwardLatencyMarker(LatencyMarker marker) {
  // all operators are tracking latencies
  this.latencyGauge.reportLatency(marker, false);
  // everything except sinks forwards latency markers
  this.output.emitLatencyMarker(marker);
}

代码示例来源:origin: org.apache.flink/flink-streaming-java_2.11

@Override
  public void onProcessingTime(long timestamp) throws Exception {
    try {
      // ProcessingTimeService callbacks are executed under the checkpointing lock
      output.emitLatencyMarker(new LatencyMarker(timestamp, operatorId, subtaskIndex));
    } catch (Throwable t) {
      // we catch the Throwables here so that we don't trigger the processing
      // timer services async exception handler
      LOG.warn("Error while emitting latency marker.", t);
    }
  }
},

相关文章