Skip to content

Commit eebfb46

Browse files
anton-kutuzovAnton Kutuzov
authored andcommitted
Add supporting of timestamp type in statistics for hive 3 and 4 versions
1 parent b8997c2 commit eebfb46

6 files changed

Lines changed: 265 additions & 45 deletions

File tree

plugin/trino-hive/src/main/java/io/trino/plugin/hive/metastore/thrift/ThriftMetastoreUtil.java

Lines changed: 44 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -43,6 +43,8 @@
4343
import io.trino.hive.thrift.metastore.SerDeInfo;
4444
import io.trino.hive.thrift.metastore.StorageDescriptor;
4545
import io.trino.hive.thrift.metastore.StringColumnStatsData;
46+
import io.trino.hive.thrift.metastore.Timestamp;
47+
import io.trino.hive.thrift.metastore.TimestampColumnStatsData;
4648
import io.trino.metastore.AcidOperation;
4749
import io.trino.metastore.Column;
4850
import io.trino.metastore.Database;
@@ -83,6 +85,8 @@
8385
import java.math.BigInteger;
8486
import java.nio.ByteBuffer;
8587
import java.time.LocalDate;
88+
import java.time.LocalDateTime;
89+
import java.time.ZoneOffset;
8690
import java.util.ArrayDeque;
8791
import java.util.Arrays;
8892
import java.util.Collection;
@@ -149,8 +153,13 @@
149153
import static io.trino.spi.type.IntegerType.INTEGER;
150154
import static io.trino.spi.type.RealType.REAL;
151155
import static io.trino.spi.type.SmallintType.SMALLINT;
156+
import static io.trino.spi.type.StandardTypes.TIMESTAMP;
157+
import static io.trino.spi.type.Timestamps.MICROSECONDS_PER_MILLISECOND;
158+
import static io.trino.spi.type.Timestamps.MICROSECONDS_PER_SECOND;
152159
import static io.trino.spi.type.TinyintType.TINYINT;
153160
import static io.trino.spi.type.VarbinaryType.VARBINARY;
161+
import static java.lang.Math.ceilDiv;
162+
import static java.lang.Math.floorDiv;
154163
import static java.lang.Math.toIntExact;
155164
import static java.lang.String.format;
156165
import static java.util.Locale.ENGLISH;
@@ -522,6 +531,10 @@ public static HiveColumnStatistics fromMetastoreApiColumnStatistics(ColumnStatis
522531
LongColumnStatsData longStatsData = columnStatistics.getStatsData().getLongStats();
523532
OptionalLong min = longStatsData.isSetLowValue() ? OptionalLong.of(longStatsData.getLowValue()) : OptionalLong.empty();
524533
OptionalLong max = longStatsData.isSetHighValue() ? OptionalLong.of(longStatsData.getHighValue()) : OptionalLong.empty();
534+
if (min.isPresent() && max.isPresent() && columnStatistics.getColType().equals(TIMESTAMP)) {
535+
min = OptionalLong.of(min.getAsLong() * MICROSECONDS_PER_SECOND);
536+
max = OptionalLong.of(max.getAsLong() * MICROSECONDS_PER_SECOND);
537+
}
525538
OptionalLong nullsCount = longStatsData.isSetNumNulls() ? fromMetastoreNullsCount(longStatsData.getNumNulls()) : OptionalLong.empty();
526539
OptionalLong distinctValuesWithNullCount = longStatsData.isSetNumDVs() ? OptionalLong.of(longStatsData.getNumDVs()) : OptionalLong.empty();
527540
return createIntegerColumnStatistics(min, max, nullsCount, distinctValuesWithNullCount);
@@ -586,6 +599,14 @@ public static HiveColumnStatistics fromMetastoreApiColumnStatistics(ColumnStatis
586599
averageColumnLength,
587600
nullsCount);
588601
}
602+
if (columnStatistics.getStatsData().isSetTimestampStats()) {
603+
TimestampColumnStatsData timestampStatsData = columnStatistics.getStatsData().getTimestampStats();
604+
OptionalLong min = timestampStatsData.isSetLowValue() ? fromMetastoreTimestamp(timestampStatsData.getLowValue()) : OptionalLong.empty();
605+
OptionalLong max = timestampStatsData.isSetHighValue() ? fromMetastoreTimestamp(timestampStatsData.getHighValue()) : OptionalLong.empty();
606+
OptionalLong nullsCount = timestampStatsData.isSetNumNulls() ? fromMetastoreNullsCount(timestampStatsData.getNumNulls()) : OptionalLong.empty();
607+
OptionalLong distinctValuesWithNullCount = timestampStatsData.isSetNumDVs() ? OptionalLong.of(timestampStatsData.getNumDVs()) : OptionalLong.empty();
608+
return createIntegerColumnStatistics(min, max, nullsCount, distinctValuesWithNullCount);
609+
}
589610
throw new TrinoException(HIVE_INVALID_METADATA, "Invalid column statistics data: " + columnStatistics);
590611
}
591612

@@ -610,6 +631,14 @@ public static OptionalLong fromMetastoreNullsCount(long nullsCount)
610631
return OptionalLong.of(nullsCount);
611632
}
612633

634+
private static OptionalLong fromMetastoreTimestamp(Timestamp timestamp)
635+
{
636+
if (timestamp == null) {
637+
return OptionalLong.empty();
638+
}
639+
return OptionalLong.of(LocalDateTime.ofEpochSecond(timestamp.getSecondsSinceEpoch(), 0, ZoneOffset.UTC).toInstant(ZoneOffset.UTC).toEpochMilli() * MICROSECONDS_PER_MILLISECOND);
640+
}
641+
613642
private static Optional<BigDecimal> fromMetastoreDecimal(@Nullable Decimal decimal)
614643
{
615644
if (decimal == null) {
@@ -769,8 +798,9 @@ public static ColumnStatisticsObj createMetastoreColumnStatistics(String columnN
769798
case SHORT:
770799
case INT:
771800
case LONG:
772-
case TIMESTAMP:
773801
return createLongStatistics(columnName, columnType, statistics);
802+
case TIMESTAMP:
803+
return createTimestampStatistics(columnName, columnType, statistics);
774804
case FLOAT:
775805
case DOUBLE:
776806
return createDoubleStatistics(columnName, columnType, statistics);
@@ -820,6 +850,18 @@ private static ColumnStatisticsObj createLongStatistics(String columnName, HiveT
820850
return new ColumnStatisticsObj(columnName, columnType.toString(), longStats(data));
821851
}
822852

853+
private static ColumnStatisticsObj createTimestampStatistics(String columnName, HiveType columnType, HiveColumnStatistics statistics)
854+
{
855+
LongColumnStatsData data = new LongColumnStatsData();
856+
statistics.getIntegerStatistics().ifPresent(timestampStatistics -> {
857+
timestampStatistics.getMin().ifPresent(value -> data.setLowValue(floorDiv(value, MICROSECONDS_PER_SECOND)));
858+
timestampStatistics.getMax().ifPresent(value -> data.setHighValue(ceilDiv(value, MICROSECONDS_PER_SECOND)));
859+
});
860+
statistics.getNullsCount().ifPresent(data::setNumNulls);
861+
statistics.getDistinctValuesWithNullCount().ifPresent(data::setNumDVs);
862+
return new ColumnStatisticsObj(columnName, columnType.toString(), longStats(data));
863+
}
864+
823865
private static ColumnStatisticsObj createDoubleStatistics(String columnName, HiveType columnType, HiveColumnStatistics statistics)
824866
{
825867
DoubleColumnStatsData data = new DoubleColumnStatsData();
@@ -894,8 +936,7 @@ public static Set<HiveColumnStatisticType> getSupportedColumnStatistics(Type typ
894936
return ImmutableSet.of(MIN_VALUE, MAX_VALUE, NUMBER_OF_DISTINCT_VALUES, NUMBER_OF_NON_NULL_VALUES);
895937
}
896938
if (type instanceof TimestampType || type instanceof TimestampWithTimeZoneType) {
897-
// TODO (https://github.com/trinodb/trino/issues/5859) Add support for timestamp MIN_VALUE, MAX_VALUE
898-
return ImmutableSet.of(NUMBER_OF_DISTINCT_VALUES, NUMBER_OF_NON_NULL_VALUES);
939+
return ImmutableSet.of(MIN_VALUE, MAX_VALUE, NUMBER_OF_DISTINCT_VALUES, NUMBER_OF_NON_NULL_VALUES);
899940
}
900941
if (type instanceof VarcharType || type instanceof CharType) {
901942
// TODO Collect MIN,MAX once it is used by the optimizer

plugin/trino-hive/src/main/java/io/trino/plugin/hive/statistics/AbstractHiveStatisticsProvider.java

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -44,6 +44,7 @@
4444
import io.trino.spi.statistics.TableStatistics;
4545
import io.trino.spi.type.CharType;
4646
import io.trino.spi.type.DecimalType;
47+
import io.trino.spi.type.TimestampType;
4748
import io.trino.spi.type.Type;
4849
import io.trino.spi.type.VarcharType;
4950

@@ -844,6 +845,9 @@ private static Optional<DoubleRange> createRange(Type type, HiveColumnStatistics
844845
if (type.equals(DATE)) {
845846
return statistics.getDateStatistics().flatMap(AbstractHiveStatisticsProvider::createDateRange);
846847
}
848+
if (type instanceof TimestampType) {
849+
return statistics.getIntegerStatistics().flatMap(AbstractHiveStatisticsProvider::createTimestampRange);
850+
}
847851
if (type instanceof DecimalType) {
848852
return statistics.getDecimalStatistics().flatMap(AbstractHiveStatisticsProvider::createDecimalRange);
849853
}
@@ -896,6 +900,14 @@ private static Optional<DoubleRange> createDateRange(DateStatistics statistics)
896900
return Optional.empty();
897901
}
898902

903+
private static Optional<DoubleRange> createTimestampRange(IntegerStatistics statistics)
904+
{
905+
if (statistics.getMin().isPresent() && statistics.getMax().isPresent()) {
906+
return Optional.of(new DoubleRange(statistics.getMin().getAsLong(), statistics.getMax().getAsLong()));
907+
}
908+
return Optional.empty();
909+
}
910+
899911
private static Optional<DoubleRange> createDecimalRange(DecimalStatistics statistics)
900912
{
901913
if (statistics.getMin().isPresent() && statistics.getMax().isPresent()) {

plugin/trino-hive/src/main/java/io/trino/plugin/hive/util/Statistics.java

Lines changed: 23 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -26,10 +26,13 @@
2626
import io.trino.spi.Page;
2727
import io.trino.spi.TrinoException;
2828
import io.trino.spi.block.Block;
29+
import io.trino.spi.block.Fixed12Block;
2930
import io.trino.spi.statistics.ColumnStatisticMetadata;
3031
import io.trino.spi.statistics.ComputedStatistics;
3132
import io.trino.spi.type.DecimalType;
3233
import io.trino.spi.type.Decimals;
34+
import io.trino.spi.type.LongTimestamp;
35+
import io.trino.spi.type.TimestampType;
3336
import io.trino.spi.type.Type;
3437

3538
import java.math.BigDecimal;
@@ -65,6 +68,7 @@
6568
import static io.trino.spi.type.IntegerType.INTEGER;
6669
import static io.trino.spi.type.RealType.REAL;
6770
import static io.trino.spi.type.SmallintType.SMALLINT;
71+
import static io.trino.spi.type.TimestampType.MAX_PRECISION;
6872
import static io.trino.spi.type.TinyintType.TINYINT;
6973
import static java.lang.Float.intBitsToFloat;
7074
import static java.lang.Math.toIntExact;
@@ -128,10 +132,12 @@ else if (type.equals(DOUBLE) || type.equals(REAL)) {
128132
else if (type.equals(DATE)) {
129133
result.setDateStatistics(new DateStatistics(Optional.empty(), Optional.empty()));
130134
}
135+
else if (type instanceof TimestampType) {
136+
result.setIntegerStatistics(new IntegerStatistics(OptionalLong.empty(), OptionalLong.empty()));
137+
}
131138
else if (type instanceof DecimalType) {
132139
result.setDecimalStatistics(new DecimalStatistics(Optional.empty(), Optional.empty()));
133140
}
134-
// TODO (https://github.com/trinodb/trino/issues/5859) Add support for timestamp
135141
else {
136142
throw new IllegalArgumentException("Unexpected type: " + type);
137143
}
@@ -240,10 +246,12 @@ else if (type.equals(DOUBLE) || type.equals(REAL)) {
240246
else if (type.equals(DATE)) {
241247
result.setDateStatistics(new DateStatistics(getDateValue(type, min), getDateValue(type, max)));
242248
}
249+
else if (type instanceof TimestampType) {
250+
result.setIntegerStatistics(new IntegerStatistics(getTimestampValue(type, min), getTimestampValue(type, max)));
251+
}
243252
else if (type instanceof DecimalType) {
244253
result.setDecimalStatistics(new DecimalStatistics(getDecimalValue(type, min), getDecimalValue(type, max)));
245254
}
246-
// TODO (https://github.com/trinodb/trino/issues/5859) Add support for timestamp
247255
else {
248256
throw new IllegalArgumentException("Unexpected type: " + type);
249257
}
@@ -288,6 +296,19 @@ private static Optional<LocalDate> getDateValue(Type type, Block block)
288296
return Optional.of(LocalDate.ofEpochDay(days));
289297
}
290298

299+
private static OptionalLong getTimestampValue(Type type, Block block)
300+
{
301+
verify(type instanceof TimestampType, "Unsupported type: %s", type);
302+
if (block.isNull(0)) {
303+
return OptionalLong.empty();
304+
}
305+
if (block instanceof Fixed12Block) {
306+
LongTimestamp ts = (LongTimestamp) TimestampType.createTimestampType(MAX_PRECISION).getObject(block, 0);
307+
return OptionalLong.of(ts.getEpochMicros());
308+
}
309+
return OptionalLong.of(type.getLong(block, 0));
310+
}
311+
291312
private static Optional<BigDecimal> getDecimalValue(Type type, Block block)
292313
{
293314
verify(type instanceof DecimalType, "Unsupported type: %s", type);

0 commit comments

Comments
 (0)