优化监控指标名,增加文件大小分布监控
This commit is contained in:
@@ -45,7 +45,7 @@ public class FileChunkCombiner {
|
||||
.filter(new FileChunkFilterFunction(configuration.getLong(Configs.FILE_MAX_SIZE), configuration.getString(Configs.FILTER_EXPRESSION)))
|
||||
.assignTimestampsAndWatermarks(watermarkStrategy);
|
||||
|
||||
OutputTag<FileChunk> delayedChunkOutputTag = new OutputTag<>("delayed-chunk") {
|
||||
OutputTag<FileChunk> delayedChunkOutputTag = new OutputTag<FileChunk>("delayed-chunk") {
|
||||
};
|
||||
|
||||
List<Trigger<Object, TimeWindow>> triggers = new ArrayList<>();
|
||||
@@ -69,7 +69,7 @@ public class FileChunkCombiner {
|
||||
windowStream.getSideOutput(delayedChunkOutputTag)
|
||||
.map(new SideOutputMapFunction())
|
||||
.addSink(new HosSink(configuration))
|
||||
.name("Hos Delayed Chunk");
|
||||
.name("Delayed Chunk");
|
||||
} else {
|
||||
windowStream.addSink(new HBaseSink(configuration))
|
||||
.name("HBase")
|
||||
@@ -77,7 +77,7 @@ public class FileChunkCombiner {
|
||||
windowStream.getSideOutput(delayedChunkOutputTag)
|
||||
.map(new SideOutputMapFunction())
|
||||
.addSink(new HBaseSink(configuration))
|
||||
.name("HBase Delayed Chunk");
|
||||
.name("Delayed Chunk");
|
||||
}
|
||||
|
||||
environment.execute(configuration.get(Configs.FLINK_JOB_NAME));
|
||||
|
||||
Reference in New Issue
Block a user