apache / apache/arrow-java

[Java] Memory leak when use VectorSchemaSlice slice

Aperta
#144 4 commenti 0 reazioni 0 assegnatari Vedi su GitHub
Lingua principale
Java
Stelle
94
Fork
152
Merge medio
3g 16h
PR unite (30g)
11

Descrizione

### Describe the usage question you have. Please include as many useful details as possible.

In my Flink process function, I receive serialized VectorSchemaRoot data which needs to be deserialized for further processing. As I need to operate on it row by row, I utilized the slice function. However, this approach can lead to an increase in direct memory usage which in turn reduces the heap memory space in Flink. Ultimately, this can result in an Out Of Memory (OOM) exception as the heap memory space becomes insufficient. However, if I serialize the sliced VectorSchemaRoot and pass it as a parameter to the downstream function, which then deserializes it again, there will be no more OOM issues.

```
public void processElement(I value, ProcessFunction.Context context, Collector collector) throws Exception {
if (this.writerHelper == null) {
initWriterHelper();
}
reader = new ArrowStreamReader(new ByteArrayInputStream((byte[]) value), ArrowUtil.rootAllocator);
try {
while (reader.loadNextBatch()) {
VectorSchemaRoot vsr = reader.getVectorSchemaRoot();
int rowCount = vsr.getRowCount();
for (int i = 0; i < rowCount; i++) {
//split to row
VectorSchemaRoot row = vsr.slice(i, 1);
// this.writerHelper.write(row); This approach can result in a memory leak, whereas the following method will not.

ByteArrayOutputStream out = new ByteArrayOutputStream();
ArrowStreamWriter writer =
new ArrowStreamWriter(row, null, Channels.newChannel(out));
writer.start();
writer.writeBatch();
this.writerHelper.write(out.toByteArray());
row.clear();
row.close();
}
vsr.clear();
vsr.close();
}
} catch (Exception ex) {
ex.printStackTrace();
} finally {
reader.close();
}
}```

### Component(s)

Java

Guida per i contributori

Apri la guida per i contributori

Direzione di ricerca

Inizia con l’esempio di processElement, in particolare con VectorSchemaRoot.slice(i, 1) e le chiamate row.clear()/row.close(). Riproduci la crescita della memoria diretta durante la scrittura di righe suddivise tramite writerHelper, quindi confrontala con il percorso di serializzazione e deserializzazione. Il lavoro è completato quando il percorso delle righe suddivise non causa più la crescita di memoria segnalata né OOM.

Scritto dal modello di indicizzazione a partire dal testo della issue.

Valutazione

Stack tecnologico
java
Ambito
data
Tipo di issue
Bug
Difficoltà
4/5
Tempo stimato
3-5 giorni
Stato di attività
Ferma
Chiarezza
Da chiarire
Idoneità per principianti
30/100

Ricevi le nuove issue nella tua casella

Un breve riepilogo di issue GitHub adatte ai principianti.