Skip to content

Commit 314ac78

Browse files
authored
[kudu] pick kudu connector to 452 version. (#232)
Migrate the changes of trino Kudu connector from version 435 to version 452. Changed some code to adapt to the changes of Trino Spi. Performed tpch tests on both Trino and Doris. It is recommended to compile with mvn package -Dmaven.test.skip=true.
1 parent 7fee875 commit 314ac78

34 files changed

Lines changed: 1007 additions & 734 deletions

plugin/trino-kudu/pom.xml

Lines changed: 27 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -10,11 +10,12 @@
1010

1111
<artifactId>trino-kudu</artifactId>
1212
<packaging>trino-plugin</packaging>
13-
<description>Trino - Kudu Connector</description>
13+
<description>Trino - Kudu connector</description>
1414

1515
<properties>
1616
<air.main.basedir>${project.parent.basedir}</air.main.basedir>
17-
<kudu.version>1.15.0</kudu.version>
17+
<air.compiler.fail-warnings>true</air.compiler.fail-warnings>
18+
<kudu.version>1.17.0</kudu.version>
1819
</properties>
1920

2021
<dependencies>
@@ -131,12 +132,30 @@
131132
<scope>provided</scope>
132133
</dependency>
133134

135+
<dependency>
136+
<groupId>com.google.errorprone</groupId>
137+
<artifactId>error_prone_annotations</artifactId>
138+
<scope>runtime</scope>
139+
</dependency>
140+
134141
<dependency>
135142
<groupId>io.airlift</groupId>
136143
<artifactId>log-manager</artifactId>
137144
<scope>runtime</scope>
138145
</dependency>
139146

147+
<dependency>
148+
<groupId>com.github.docker-java</groupId>
149+
<artifactId>docker-java-api</artifactId>
150+
<scope>test</scope>
151+
</dependency>
152+
153+
<dependency>
154+
<groupId>dev.failsafe</groupId>
155+
<artifactId>failsafe</artifactId>
156+
<scope>test</scope>
157+
</dependency>
158+
140159
<dependency>
141160
<groupId>io.airlift</groupId>
142161
<artifactId>junit-extensions</artifactId>
@@ -218,20 +237,21 @@
218237
</dependency>
219238

220239
<dependency>
221-
<groupId>org.testcontainers</groupId>
222-
<artifactId>testcontainers</artifactId>
240+
<groupId>org.rnorth.duct-tape</groupId>
241+
<artifactId>duct-tape</artifactId>
242+
<version>1.0.8</version>
223243
<scope>test</scope>
224244
</dependency>
225245

226246
<dependency>
227247
<groupId>org.testcontainers</groupId>
228-
<artifactId>toxiproxy</artifactId>
248+
<artifactId>testcontainers</artifactId>
229249
<scope>test</scope>
230250
</dependency>
231251

232252
<dependency>
233-
<groupId>org.testng</groupId>
234-
<artifactId>testng</artifactId>
253+
<groupId>org.testcontainers</groupId>
254+
<artifactId>toxiproxy</artifactId>
235255
<scope>test</scope>
236256
</dependency>
237257
</dependencies>

plugin/trino-kudu/src/main/java/io/trino/plugin/kudu/KuduClientConfig.java

Lines changed: 5 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -13,7 +13,6 @@
1313
*/
1414
package io.trino.plugin.kudu;
1515

16-
import com.google.common.base.Splitter;
1716
import com.google.common.collect.ImmutableList;
1817
import io.airlift.configuration.Config;
1918
import io.airlift.configuration.ConfigDescription;
@@ -35,11 +34,11 @@
3534
@DefunctConfig("kudu.client.default-socket-read-timeout")
3635
public class KuduClientConfig
3736
{
38-
private static final Splitter SPLITTER = Splitter.on(',').trimResults().omitEmptyStrings();
37+
private static final Duration DEFAULT_OPERATION_TIMEOUT = new Duration(30, TimeUnit.SECONDS);
3938

4039
private List<String> masterAddresses = ImmutableList.of();
41-
private Duration defaultAdminOperationTimeout = new Duration(30, TimeUnit.SECONDS);
42-
private Duration defaultOperationTimeout = new Duration(30, TimeUnit.SECONDS);
40+
private Duration defaultAdminOperationTimeout = DEFAULT_OPERATION_TIMEOUT;
41+
private Duration defaultOperationTimeout = DEFAULT_OPERATION_TIMEOUT;
4342
private boolean disableStatistics;
4443
private boolean schemaEmulationEnabled;
4544
private String schemaEmulationPrefix = "presto::";
@@ -54,15 +53,9 @@ public List<String> getMasterAddresses()
5453
}
5554

5655
@Config("kudu.client.master-addresses")
57-
public KuduClientConfig setMasterAddresses(String commaSeparatedList)
56+
public KuduClientConfig setMasterAddresses(List<String> commaSeparatedList)
5857
{
59-
this.masterAddresses = SPLITTER.splitToList(commaSeparatedList);
60-
return this;
61-
}
62-
63-
public KuduClientConfig setMasterAddresses(String... contactPoints)
64-
{
65-
this.masterAddresses = ImmutableList.copyOf(contactPoints);
58+
this.masterAddresses = commaSeparatedList;
6659
return this;
6760
}
6861

plugin/trino-kudu/src/main/java/io/trino/plugin/kudu/KuduClientSession.java

Lines changed: 44 additions & 37 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,8 @@
1515

1616
import com.google.common.collect.ImmutableList;
1717
import io.airlift.log.Logger;
18+
import io.airlift.slice.Slice;
19+
import io.airlift.slice.Slices;
1820
import io.trino.plugin.kudu.properties.ColumnDesign;
1921
import io.trino.plugin.kudu.properties.HashPartitionDefinition;
2022
import io.trino.plugin.kudu.properties.KuduTableProperties;
@@ -58,6 +60,7 @@
5860

5961
import java.io.IOException;
6062
import java.math.BigDecimal;
63+
import java.nio.ByteBuffer;
6164
import java.util.ArrayList;
6265
import java.util.List;
6366
import java.util.Map;
@@ -178,7 +181,7 @@ public List<KuduSplit> buildKuduSplits(KuduTableHandle tableHandle, DynamicFilte
178181
.boxed().collect(toList());
179182
for (ColumnHandle column : desiredColumns.get()) {
180183
KuduColumnHandle k = (KuduColumnHandle) column;
181-
int index = k.getOrdinalPosition();
184+
int index = k.ordinalPosition();
182185
if (index >= primaryKeyColumnCount) {
183186
columnIndexes.add(index);
184187
}
@@ -194,7 +197,7 @@ public List<KuduSplit> buildKuduSplits(KuduTableHandle tableHandle, DynamicFilte
194197
else {
195198
if (desiredColumns.isPresent()) {
196199
columnIndexes = desiredColumns.get().stream()
197-
.map(handle -> ((KuduColumnHandle) handle).getOrdinalPosition())
200+
.map(handle -> ((KuduColumnHandle) handle).ordinalPosition())
198201
.collect(toImmutableList());
199202
}
200203
else {
@@ -323,12 +326,12 @@ public void addColumn(SchemaTableName schemaTableName, ColumnMetadata column)
323326
String rawName = schemaEmulation.toRawName(schemaTableName);
324327
AlterTableOptions alterOptions = new AlterTableOptions();
325328
Type type = TypeHelper.toKuduClientType(column.getType());
326-
alterOptions.addColumn(
327-
new ColumnSchemaBuilder(column.getName(), type)
328-
.nullable(true)
329-
.defaultValue(null)
330-
.comment(nullToEmpty(column.getComment())) // Kudu doesn't allow null comment
331-
.build());
329+
ColumnSchemaBuilder builder = new ColumnSchemaBuilder(column.getName(), type)
330+
.nullable(true)
331+
.defaultValue(null)
332+
.comment(nullToEmpty(column.getComment())); // Kudu doesn't allow null comment
333+
setTypeAttributes(column, builder);
334+
alterOptions.addColumn(builder.build());
332335
client.alterTable(rawName, alterOptions);
333336
}
334337
catch (KuduException e) {
@@ -384,8 +387,8 @@ private void changeRangePartition(SchemaTableName schemaTableName, RangePartitio
384387
if (definition == null) {
385388
throw new TrinoException(QUERY_REJECTED, "Table " + schemaTableName + " has no range partition");
386389
}
387-
PartialRow lowerBound = KuduTableProperties.toRangeBoundToPartialRow(schema, definition, rangePartition.getLower());
388-
PartialRow upperBound = KuduTableProperties.toRangeBoundToPartialRow(schema, definition, rangePartition.getUpper());
390+
PartialRow lowerBound = KuduTableProperties.toRangeBoundToPartialRow(schema, definition, rangePartition.lower());
391+
PartialRow upperBound = KuduTableProperties.toRangeBoundToPartialRow(schema, definition, rangePartition.upper());
389392
AlterTableOptions alterOptions = new AlterTableOptions();
390393
switch (change) {
391394
case ADD:
@@ -467,19 +470,19 @@ private CreateTableOptions buildCreateTableOptions(Schema schema, Map<String, Ob
467470
PartitionDesign partitionDesign = KuduTableProperties.getPartitionDesign(properties);
468471
if (partitionDesign.getHash() != null) {
469472
for (HashPartitionDefinition partition : partitionDesign.getHash()) {
470-
options.addHashPartitions(partition.getColumns(), partition.getBuckets());
473+
options.addHashPartitions(partition.columns(), partition.buckets());
471474
}
472475
}
473476
if (partitionDesign.getRange() != null) {
474477
rangePartitionDefinition = partitionDesign.getRange();
475-
options.setRangePartitionColumns(rangePartitionDefinition.getColumns());
478+
options.setRangePartitionColumns(rangePartitionDefinition.columns());
476479
}
477480

478481
List<RangePartition> rangePartitions = KuduTableProperties.getRangePartitions(properties);
479482
if (rangePartitionDefinition != null && !rangePartitions.isEmpty()) {
480483
for (RangePartition rangePartition : rangePartitions) {
481-
PartialRow lower = KuduTableProperties.toRangeBoundToPartialRow(schema, rangePartitionDefinition, rangePartition.getLower());
482-
PartialRow upper = KuduTableProperties.toRangeBoundToPartialRow(schema, rangePartitionDefinition, rangePartition.getUpper());
484+
PartialRow lower = KuduTableProperties.toRangeBoundToPartialRow(schema, rangePartitionDefinition, rangePartition.lower());
485+
PartialRow upper = KuduTableProperties.toRangeBoundToPartialRow(schema, rangePartitionDefinition, rangePartition.upper());
483486
options.addRangePartition(lower, upper);
484487
}
485488
}
@@ -503,7 +506,7 @@ private void addConstraintPredicates(KuduTable table, KuduScanToken.KuduScanToke
503506

504507
Schema schema = table.getSchema();
505508
constraintSummary.getDomains().orElseThrow().forEach((columnHandle, domain) -> {
506-
int position = ((KuduColumnHandle) columnHandle).getOrdinalPosition();
509+
int position = ((KuduColumnHandle) columnHandle).ordinalPosition();
507510
ColumnSchema columnSchema = schema.getColumnByIndex(position);
508511
verify(!domain.isNone(), "Domain is none");
509512
if (domain.isAll()) {
@@ -529,8 +532,8 @@ else if (domain.isSingleValue()) {
529532
KuduPredicate predicate = createInListPredicate(columnSchema, discreteValues);
530533
builder.addPredicate(predicate);
531534
}
532-
else if (valueSet instanceof SortedRangeSet) {
533-
Ranges ranges = ((SortedRangeSet) valueSet).getRanges();
535+
else if (valueSet instanceof SortedRangeSet sortedRangeSet) {
536+
Ranges ranges = sortedRangeSet.getRanges();
534537
List<Range> rangeList = ranges.getOrderedRanges();
535538
if (rangeList.stream().allMatch(Range::isSingleValue)) {
536539
io.trino.spi.type.Type type = TypeHelper.fromKuduColumn(columnSchema);
@@ -577,35 +580,39 @@ private KuduPredicate createComparisonPredicate(ColumnSchema columnSchema, KuduP
577580
{
578581
io.trino.spi.type.Type type = TypeHelper.fromKuduColumn(columnSchema);
579582
Object javaValue = TypeHelper.getJavaValue(type, value);
580-
if (javaValue instanceof Long) {
581-
return KuduPredicate.newComparisonPredicate(columnSchema, op, (Long) javaValue);
583+
if (javaValue instanceof Long longValue) {
584+
return KuduPredicate.newComparisonPredicate(columnSchema, op, longValue);
582585
}
583-
if (javaValue instanceof BigDecimal) {
584-
return KuduPredicate.newComparisonPredicate(columnSchema, op, (BigDecimal) javaValue);
586+
if (javaValue instanceof BigDecimal bigDecimal) {
587+
return KuduPredicate.newComparisonPredicate(columnSchema, op, bigDecimal);
585588
}
586-
if (javaValue instanceof Integer) {
587-
return KuduPredicate.newComparisonPredicate(columnSchema, op, (Integer) javaValue);
589+
if (javaValue instanceof Integer integerValue) {
590+
return KuduPredicate.newComparisonPredicate(columnSchema, op, integerValue);
588591
}
589-
if (javaValue instanceof Short) {
590-
return KuduPredicate.newComparisonPredicate(columnSchema, op, (Short) javaValue);
592+
if (javaValue instanceof Short shortValue) {
593+
return KuduPredicate.newComparisonPredicate(columnSchema, op, shortValue);
591594
}
592-
if (javaValue instanceof Byte) {
593-
return KuduPredicate.newComparisonPredicate(columnSchema, op, (Byte) javaValue);
595+
if (javaValue instanceof Byte byteValue) {
596+
return KuduPredicate.newComparisonPredicate(columnSchema, op, byteValue);
594597
}
595-
if (javaValue instanceof String) {
596-
return KuduPredicate.newComparisonPredicate(columnSchema, op, (String) javaValue);
598+
if (javaValue instanceof String stringValue) {
599+
return KuduPredicate.newComparisonPredicate(columnSchema, op, stringValue);
597600
}
598-
if (javaValue instanceof Double) {
599-
return KuduPredicate.newComparisonPredicate(columnSchema, op, (Double) javaValue);
601+
if (javaValue instanceof Double doubleValue) {
602+
return KuduPredicate.newComparisonPredicate(columnSchema, op, doubleValue);
600603
}
601-
if (javaValue instanceof Float) {
602-
return KuduPredicate.newComparisonPredicate(columnSchema, op, (Float) javaValue);
604+
if (javaValue instanceof Float floatValue) {
605+
return KuduPredicate.newComparisonPredicate(columnSchema, op, floatValue);
603606
}
604-
if (javaValue instanceof Boolean) {
605-
return KuduPredicate.newComparisonPredicate(columnSchema, op, (Boolean) javaValue);
607+
if (javaValue instanceof Boolean booleanValue) {
608+
return KuduPredicate.newComparisonPredicate(columnSchema, op, booleanValue);
606609
}
607-
if (javaValue instanceof byte[]) {
608-
return KuduPredicate.newComparisonPredicate(columnSchema, op, (byte[]) javaValue);
610+
if (javaValue instanceof byte[] byteArrayValue) {
611+
return KuduPredicate.newComparisonPredicate(columnSchema, op, byteArrayValue);
612+
}
613+
if (javaValue instanceof ByteBuffer byteBuffer) {
614+
Slice slice = Slices.wrappedHeapBuffer(byteBuffer);
615+
return KuduPredicate.newComparisonPredicate(columnSchema, op, slice.getBytes(0, slice.length()));
609616
}
610617
if (javaValue == null) {
611618
throw new IllegalStateException("Unexpected null java value for column " + columnSchema.getName());

plugin/trino-kudu/src/main/java/io/trino/plugin/kudu/KuduColumnHandle.java

Lines changed: 5 additions & 71 deletions
Original file line numberDiff line numberDiff line change
@@ -13,60 +13,28 @@
1313
*/
1414
package io.trino.plugin.kudu;
1515

16-
import com.fasterxml.jackson.annotation.JsonCreator;
17-
import com.fasterxml.jackson.annotation.JsonProperty;
1816
import io.trino.spi.connector.ColumnHandle;
1917
import io.trino.spi.connector.ColumnMetadata;
2018
import io.trino.spi.type.Type;
2119
import io.trino.spi.type.VarbinaryType;
2220

23-
import java.util.Objects;
24-
25-
import static com.google.common.base.MoreObjects.toStringHelper;
2621
import static java.util.Objects.requireNonNull;
2722

28-
public class KuduColumnHandle
23+
public record KuduColumnHandle(String name, int ordinalPosition, Type type)
2924
implements ColumnHandle
3025
{
3126
public static final String ROW_ID = "row_uuid";
3227
public static final int ROW_ID_POSITION = -1;
3328

3429
public static final KuduColumnHandle ROW_ID_HANDLE = new KuduColumnHandle(ROW_ID, ROW_ID_POSITION, VarbinaryType.VARBINARY);
3530

36-
private final String name;
37-
private final int ordinalPosition;
38-
private final Type type;
39-
40-
@JsonCreator
41-
public KuduColumnHandle(
42-
@JsonProperty("name") String name,
43-
@JsonProperty("ordinalPosition") int ordinalPosition,
44-
@JsonProperty("type") Type type)
45-
{
46-
this.name = requireNonNull(name, "name is null");
47-
this.ordinalPosition = ordinalPosition;
48-
this.type = requireNonNull(type, "type is null");
49-
}
50-
51-
@JsonProperty
52-
public String getName()
53-
{
54-
return name;
55-
}
56-
57-
@JsonProperty
58-
public int getOrdinalPosition()
59-
{
60-
return ordinalPosition;
61-
}
62-
63-
@JsonProperty
64-
public Type getType()
31+
public KuduColumnHandle
6532
{
66-
return type;
33+
requireNonNull(name, "name is null");
34+
requireNonNull(type, "type is null");
6735
}
6836

69-
public ColumnMetadata getColumnMetadata()
37+
public ColumnMetadata columnMetadata()
7038
{
7139
return new ColumnMetadata(name, type);
7240
}
@@ -75,38 +43,4 @@ public boolean isVirtualRowId()
7543
{
7644
return name.equals(ROW_ID);
7745
}
78-
79-
@Override
80-
public int hashCode()
81-
{
82-
return Objects.hash(
83-
name,
84-
ordinalPosition,
85-
type);
86-
}
87-
88-
@Override
89-
public boolean equals(Object obj)
90-
{
91-
if (this == obj) {
92-
return true;
93-
}
94-
if (obj == null || getClass() != obj.getClass()) {
95-
return false;
96-
}
97-
KuduColumnHandle other = (KuduColumnHandle) obj;
98-
return Objects.equals(this.name, other.name) &&
99-
Objects.equals(this.ordinalPosition, other.ordinalPosition) &&
100-
Objects.equals(this.type, other.type);
101-
}
102-
103-
@Override
104-
public String toString()
105-
{
106-
return toStringHelper(this)
107-
.add("name", name)
108-
.add("ordinalPosition", ordinalPosition)
109-
.add("type", type)
110-
.toString();
111-
}
11246
}

0 commit comments

Comments
 (0)