Масштабируемый исходный код для чтения файлов CSV в FlinkJAVA

Программисты JAVA общаются здесь
Anonymous
Масштабируемый исходный код для чтения файлов CSV в Flink

Сообщение Anonymous »

Я пытаюсь создать csvFileReader во flink, который можно масштабировать (т. е. увеличивать параллелизм), но при этом сохранять порядок. Я не уверен, связана ли моя проблема с невозможностью создать правильную стратегию водяных знаков или с тем, что я использую неправильный тип функции.
В настоящее время моя единственная реализация выглядит следующим образом: :

Код: Выделить всё

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 указана временная метка, а функциям процесса (имеющим собственное настраиваемое оконное оформление) нужны эти временные метки/события, чтобы они поступали по порядку. Так как они не будут знать, когда смыть воду из окна.
например. ввода 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

Вернуться в «JAVA»