apache / apache/arrow-java

[JAVA] Use prepared statement leads Memory leak

Đang mở
#50 0 bình luận 0 reaction 0 người được giao Xem trên GitHub
Type: bug
Ngôn ngữ chính
Java
Star
94
Fork
152
Merge trung bình
3 ngày 16 giờ
Pull request đã merge (30 ngày)
11

Mô tả

### Describe the bug, including details regarding any error messages, version, and platform.

Hi arrow maintainer, i find some memory leak in my flight-sql usage. I use arrow-16.0.0, and works on AMD Ryzen 3970X machine.
```java
package org.example;

import org.apache.arrow.flight.*;
import org.apache.arrow.flight.grpc.CredentialCallOption;
import org.apache.arrow.flight.sql.FlightSqlClient;
import org.apache.arrow.memory.BufferAllocator;
import org.apache.arrow.memory.RootAllocator;
import org.apache.arrow.vector.*;
import org.apache.arrow.vector.types.pojo.Field;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.util.*;

public class SqlRunner {

private static final Logger log = LoggerFactory.getLogger(SqlRunner.class);

static void run_flight_sql() throws Exception {
try (BufferAllocator allocator = new RootAllocator(Integer.MAX_VALUE)) {
final Location clientLocation = Location.forGrpcInsecure("127.0.0.1", 8360);
try (FlightClient client = FlightClient.builder(allocator, clientLocation).build();
FlightSqlClient sqlClient = new FlightSqlClient(client)) {

Optional credentialCallOption = client.authenticateBasicToken("admin", "public");
CallHeaders headers = new FlightCallHeaders();
headers.insert("database", "test");

Set options = new HashSet<>();
credentialCallOption.ifPresent(options::add);
options.add(new HeaderCallOption(headers));
// use the sql to query is ok
try {
String query = "SELECT count(*) from test.sx1;";
executeQuery(sqlClient, query, options);
} catch (Exception e){
e.printStackTrace();
throw e;
}

// where memory leak happens
try {
try (FlightSqlClient.PreparedStatement preparedStatement = sqlClient.prepare("insert into table sx1 (sid, value, flag) values(?, ?, ?);", options.toArray(new CallOption[0]))) {
insertBatch(sqlClient, preparedStatement, allocator, options);
}
} catch (Exception e){
e.printStackTrace();
}
}
}
}

private static void executeQuery(FlightSqlClient sqlClient, String query, Set options) throws Exception {
final FlightInfo info = sqlClient.execute(query, options.toArray(new CallOption[0]));
final Ticket ticket = info.getEndpoints().get(0).getTicket();
try (FlightStream stream = sqlClient.getStream(ticket, options.toArray(new CallOption[0]))) {
while (stream.next()) {
try (VectorSchemaRoot schemaRoot = stream.getRoot()) {
log.info(schemaRoot.contentToTSVString());
}
}
}
}

private static void insertBatch(FlightSqlClient sqlClient, FlightSqlClient.PreparedStatement preparedStatement, BufferAllocator allocator, Set options) throws Exception {
try (IntVector sids = new IntVector("sid", allocator);
Float4Vector values = new Float4Vector("value", allocator);
TinyIntVector flags = new TinyIntVector("flag", allocator)) {

sids.allocateNew(100);
values.allocateNew(100);
flags.allocateNew(100);

for (int i = 0; i < 100; i++) {
sids.setSafe(i, i);
values.setSafe(i, (float) i);
flags.setSafe(i, (byte) i);
}

List fields = Arrays.asList(sids.getField(), values.getField(), flags.getField());
List fieldVectors = Arrays.asList(sids, values, flags);
try (VectorSchemaRoot vectorSchemaRoot = new VectorSchemaRoot(fields, fieldVectors)){
vectorSchemaRoot.setRowCount(100);

preparedStatement.setParameters(vectorSchemaRoot);
FlightInfo info = preparedStatement.execute();

final Ticket ticket = info.getEndpoints().get(0).getTicket();
try (FlightStream stream = sqlClient.getStream(ticket, options.toArray(new CallOption[0]))) {
while (stream.next()) {
try (VectorSchemaRoot schemaRoot = stream.getRoot()) {
List vectors = schemaRoot.getFieldVectors();
for (int i = 0; i < vectors.size(); i++) {
System.out.printf("%d %s\n", i, vectors.get(i));
}
}
}
}
preparedStatement.clearParameters();
}
}
}

public static void main(String[] args) throws Exception {
run_flight_sql();
}
}
```

That's the memory allocator verboseString, is it a bug? or server side bad implement?
```bash
=================================================================
Allocator(ROOT) 0/16/1424/2147483647 (res/actual/peak/limit)
child allocators: 1
Allocator(flight-client) 0/16/272/9223372036854775807 (res/actual/peak/limit)
child allocators: 0
ledgers: 1
ledger[6] allocator: flight-client), isOwning: , size: , references: 1, life: 1242736490302663..0, allocatorManager: [, life: ] holds 1 buffers.
ArrowBuf[23], address:139720841494544, capacity:16
event log for: ArrowBuf[23]
1242736490544534 create()
at org.apache.arrow.memory.util.HistoricalLog$Event.(HistoricalLog.java:180)
at org.apache.arrow.memory.util.HistoricalLog.recordEvent(HistoricalLog.java:85)
at org.apache.arrow.memory.ArrowBuf.(ArrowBuf.java:98)
at org.apache.arrow.memory.BufferLedger.newArrowBuf(BufferLedger.java:259)
at org.apache.arrow.memory.BaseAllocator.bufferWithoutReservation(BaseAllocator.java:352)
at org.apache.arrow.memory.BaseAllocator.buffer(BaseAllocator.java:328)
at org.apache.arrow.memory.BaseAllocator.buffer(BaseAllocator.java:291)
at org.apache.arrow.flight.PutResult.fromProtocol(PutResult.java:82)
at org.apache.arrow.flight.FlightClient$SetStreamObserver.onNext(FlightClient.java:466)
at org.apache.arrow.flight.FlightClient$SetStreamObserver.onNext(FlightClient.java:454)
at io.grpc.stub.ClientCalls$StreamObserverToCallListenerAdapter.onMessage(ClientCalls.java:468)
at io.grpc.ForwardingClientCallListener.onMessage(ForwardingClientCallListener.java:33)
at io.grpc.internal.ClientCallImpl$ClientStreamListenerImpl$1MessagesAvailable.runInternal(ClientCallImpl.java:657)
at io.grpc.internal.ClientCallImpl$ClientStreamListenerImpl$1MessagesAvailable.runInContext(ClientCallImpl.java:644)
at io.grpc.internal.ContextRunnable.run(ContextRunnable.java:37)
at io.grpc.internal.SerializingExecutor.run(SerializingExecutor.java:133)
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
at java.lang.Thread.run(Thread.java:750)

reservations: 0
ledgers: 0
reservations: 0

=================================================================
```

### Component(s)

Java

Hướng dẫn đóng góp

Mở hướng dẫn đóng góp

Hướng nghiên cứu

Tái hiện luồng prepared statement trong ví dụ Java được cung cấp và kiểm tra đầu ra của allocator sau khi run_flight_sql hoàn tất. Bắt đầu với các stack frame của org.apache.arrow.flight.PutResult.fromProtocol và FlightClient.SetStreamObserver.onNext, sau đó so sánh các đường đi của truy vấn và prepared statement. Hoàn tất có nghĩa là xác định liệu ArrowBuf còn lại thuộc phía client hay server, đồng thời ghi lại hoặc sửa ownership của nó.

Do mô hình lập chỉ mục viết ra từ nội dung của issue.

Đánh giá

Công nghệ
java
Lĩnh vực
databases
Loại issue
Lỗi
Độ khó
3/5
Thời gian dự kiến
1-2 ngày
Mức độ hoạt động
Đình trệ
Độ rõ ràng
Khá rõ ràng
Mức phù hợp với người mới
38/100

Nhận issue mới trong hộp thư của bạn

Bản tóm tắt ngắn những issue GitHub phù hợp với người mới.