-
Notifications
You must be signed in to change notification settings - Fork 10
Expand file tree
/
Copy pathTestUtils.java
More file actions
92 lines (82 loc) · 3.8 KB
/
Copy pathTestUtils.java
File metadata and controls
92 lines (82 loc) · 3.8 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
/*
* The MIT License
*
* Permission is hereby granted, free of charge, to any person obtaining a copy
* of this software and associated documentation files (the "Software"), to deal
* in the Software without restriction, including without limitation the rights
* to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
* copies of the Software, and to permit persons to whom the Software is
* furnished to do so, subject to the following conditions:
*
* The above copyright notice and this permission notice shall be included in
* all copies or substantial portions of the Software.
*
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
* FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
* AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
* LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
* OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN
* THE SOFTWARE.
*/
package com.influxdb.v3.client;
import java.net.URI;
import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.List;
import javax.annotation.Nonnull;
import org.apache.arrow.flight.FlightServer;
import org.apache.arrow.flight.Location;
import org.apache.arrow.flight.NoOpFlightProducer;
import org.apache.arrow.flight.Ticket;
import org.apache.arrow.memory.BufferAllocator;
import org.apache.arrow.memory.RootAllocator;
import org.apache.arrow.vector.VarCharVector;
import org.apache.arrow.vector.VectorSchemaRoot;
import org.apache.arrow.vector.types.pojo.ArrowType;
import org.apache.arrow.vector.types.pojo.Field;
import org.apache.arrow.vector.types.pojo.FieldType;
import org.apache.arrow.vector.types.pojo.Schema;
public final class TestUtils {
private TestUtils() {
throw new IllegalStateException("Utility class");
}
public static FlightServer simpleFlightServer(@Nonnull final URI uri,
@Nonnull final BufferAllocator allocator,
@Nonnull final NoOpFlightProducer producer) throws Exception {
Location location = Location.forGrpcInsecure(uri.getHost(), uri.getPort());
return FlightServer.builder(allocator, location, producer).build();
}
public static NoOpFlightProducer simpleProducer(@Nonnull final VectorSchemaRoot vectorSchemaRoot) {
return new NoOpFlightProducer() {
@Override
public void getStream(final CallContext context,
final Ticket ticket,
final ServerStreamListener listener) {
listener.start(vectorSchemaRoot);
if (listener.isReady()) {
listener.putNext();
}
listener.completed();
}
};
}
public static VectorSchemaRoot generateVectorSchemaRoot(final int fieldCount, final int rowCount) {
List<Field> fields = new ArrayList<>();
for (int i = 0; i < fieldCount; i++) {
Field field = new Field("field" + i, FieldType.nullable(new ArrowType.Utf8()), null);
fields.add(field);
}
Schema schema = new Schema(fields);
VectorSchemaRoot vectorSchemaRoot = VectorSchemaRoot.create(schema, new RootAllocator(Long.MAX_VALUE));
for (Field field : fields) {
VarCharVector vector = (VarCharVector) vectorSchemaRoot.getVector(field);
vector.allocateNew(rowCount);
for (int i = 0; i < rowCount; i++) {
vector.set(i, "Value".getBytes(StandardCharsets.UTF_8));
}
}
vectorSchemaRoot.setRowCount(rowCount);
return vectorSchemaRoot;
}
}