В настоящее время моя единственная реализация выглядит следующим образом: :
Код: Выделить всё
this.watermarkStrategy = WatermarkStrategy
.forMonotonousTimestamps()
.withTimestampAssigner((element, recordTimestamp) -> element.value.timeStamp);
.
.
.
DataStream mainStream = env.readFile( new TextInputFormat(new Path(csvFilePath)), csvFilePath, FileProcessingMode.PROCESS_ONCE, 1000)
.setParallelism(1)
.flatMap(new CSVSourceFlatMapParallelized())
.setParallelism(1)
.assignTimestampsAndWatermarks(watermarkStrategy).name("source");
например. ввода CSV-файла:
Код: Выделить всё
key,val,timestamp
E,0,500
F,1,500
Y,2,500
Z,3,500
F,4,500
Код: Выделить всё
long currentTime;
long endWindowTime;
long windowTime; //in ms
long startWindowTime;
//sum + count for each key
private volatile HashMap SumCountMap;
@Override
public void processElement(EventBasic event, Context ctx, Collector out) throws Exception {
String key = event.key;
if(event.value.timeStamp > endWindowTime){
printMapState();
outputMaxValues(out);
startWindowTime = endWindowTime;
endWindowTime += windowTime;
currentTime = event.value.timeStamp;
}
if(event.value.timeStamp < startWindowTime){
System.out.println("not working");
}
// If no maximum value has been stored yet or the incoming value is greater, update the MapState
if(!SumCountMap.containsKey(key)){
SumCountMap.put(key, new Tuple2(event.value.valueInt , 1));
}
else{
Tuple2 curr = SumCountMap.get(key);
SumCountMap.put(key, new Tuple2(curr.f0 + event.value.valueInt , curr.f1 + 1));
}
}
Код: Выделить всё
DataStream mainStream = env.addSource(new CsvEventBasicSourceFunction(csvFilePath))
.setParallelism(3)
.assignTimestampsAndWatermarks(watermarkStrategy).name("source");
Подробнее здесь: https://stackoverflow.com/questions/784 ... e-in-flink