本文整理了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
暂无
代码示例来源: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);
}
}
},
内容来源于网络,如有侵权,请联系作者删除!