Fix parquet timestamp column reading

This commit is contained in:
NeerajUnnikrishnan 2021-06-18 10:52:02 -04:00
parent 70f7dc0550
commit 2d0101dc3d
9 changed files with 41 additions and 19 deletions

View File

@ -213,6 +213,7 @@ public class ParquetPageSourceFactory
messageColumnIO,
blocks.build(),
dataSource,
timeZone,
systemMemoryContext,
maxReadBlockSize);

View File

@ -96,7 +96,6 @@ import static io.prestosql.plugin.hive.HiveTestUtils.SESSION;
import static io.prestosql.plugin.hive.HiveTestUtils.TYPE_MANAGER;
import static io.prestosql.plugin.hive.HiveTestUtils.isDistinctFrom;
import static io.prestosql.plugin.hive.HiveTestUtils.mapType;
import static io.prestosql.plugin.hive.HiveType.HIVE_TIMESTAMP;
import static io.prestosql.plugin.hive.HiveUtil.isStructuralType;
import static io.prestosql.plugin.hive.util.SerDeUtils.serializeObject;
import static io.prestosql.spi.type.BigintType.BIGINT;

View File

@ -447,8 +447,10 @@ public class TestHiveFileFormats
// TODO: empty arrays or maps with null keys don't seem to work
// Parquet does not support DATE
return TEST_COLUMNS.stream()
.filter(column -> !ImmutableSet.of("t_null_array_int", "t_array_empty", "t_map_null_key", "t_map_null_key_complex_value", "t_map_null_key_complex_key_value")
.contains(column.getName()))
.filter(TestHiveFileFormats::withoutTimestamps)
.filter(TestHiveFileFormats::withoutNullMapKeyTests)
.filter(column -> !column.getName().equals("t_null_array_int"))
.filter(column -> !column.getName().equals("t_array_empty"))
.filter(column -> column.isPartitionKey() || (
!hasType(column.getObjectInspector(), PrimitiveCategory.DATE)) &&
!hasType(column.getObjectInspector(), PrimitiveCategory.SHORT) &&

View File

@ -166,7 +166,7 @@ public enum FileFormat
@Override
public ConnectorPageSource createFileFormatReader(ConnectorSession session, HdfsEnvironment hdfsEnvironment, File targetFile, List<String> columnNames, List<Type> columnTypes)
{
HivePageSourceFactory pageSourceFactory = new ParquetPageSourceFactory(TYPE_MANAGER, hdfsEnvironment, new FileFormatDataSourceStats(), new HiveConfig());
HivePageSourceFactory pageSourceFactory = new ParquetPageSourceFactory(TYPE_MANAGER, hdfsEnvironment, new FileFormatDataSourceStats(), new HiveConfig().setParquetTimeZone("UTC"));
return createPageSource(pageSourceFactory, session, targetFile, columnNames, columnTypes, HiveStorageFormat.PARQUET);
}
@ -246,7 +246,7 @@ public enum FileFormat
@Override
public ConnectorPageSource createFileFormatReader(ConnectorSession session, HdfsEnvironment hdfsEnvironment, File targetFile, List<String> columnNames, List<Type> columnTypes)
{
HivePageSourceFactory pageSourceFactory = new ParquetPageSourceFactory(TYPE_MANAGER, hdfsEnvironment, new FileFormatDataSourceStats(), new HiveConfig());
HivePageSourceFactory pageSourceFactory = new ParquetPageSourceFactory(TYPE_MANAGER, hdfsEnvironment, new FileFormatDataSourceStats(), new HiveConfig().setParquetTimeZone("UTC"));
return createPageSource(pageSourceFactory, session, targetFile, columnNames, columnTypes, HiveStorageFormat.PARQUET);
}

View File

@ -41,7 +41,6 @@ import org.testng.annotations.Test;
import java.math.BigDecimal;
import java.math.BigInteger;
import java.time.LocalDate;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
@ -1772,7 +1771,7 @@ public abstract class AbstractTestParquetReader
if (input == null) {
return null;
}
Timestamp timestamp = new Timestamp();
long seconds = (input / 1000);
int nanos = ((input % 1000) * 1_000_000);
@ -1787,9 +1786,7 @@ public abstract class AbstractTestParquetReader
nanos -= 1_000_000_000;
seconds += 1;
}
timestamp.setTimeInMillis(seconds * 1000);
timestamp.setNanos(nanos);
return timestamp;
return Timestamp.ofEpochSecond(seconds, nanos);
}
private static SqlTimestamp intToSqlTimestamp(Integer input)
@ -1797,7 +1794,7 @@ public abstract class AbstractTestParquetReader
if (input == null) {
return null;
}
return sqlTimestampOf(input);
return sqlTimestampOf((long) input);
}
private static Date intToDate(Integer input)
@ -1805,7 +1802,7 @@ public abstract class AbstractTestParquetReader
if (input == null) {
return null;
}
return Date.valueOf(LocalDate.ofEpochDay(input).toString());
return Date.ofEpochDay(input);
}
private static SqlDate intToSqlDate(Integer input)

View File

@ -52,6 +52,11 @@
<artifactId>fastutil</artifactId>
</dependency>
<dependency>
<groupId>joda-time</groupId>
<artifactId>joda-time</artifactId>
</dependency>
<dependency>
<groupId>org.xerial.snappy</groupId>
<artifactId>snappy-java</artifactId>

View File

@ -39,6 +39,7 @@ import org.apache.parquet.hadoop.metadata.ColumnChunkMetaData;
import org.apache.parquet.hadoop.metadata.ColumnPath;
import org.apache.parquet.io.MessageColumnIO;
import org.apache.parquet.io.PrimitiveColumnIO;
import org.joda.time.DateTimeZone;
import java.io.Closeable;
import java.io.IOException;
@ -67,6 +68,7 @@ public class ParquetReader
private final List<BlockMetaData> blocks;
private final List<PrimitiveColumnIO> columns;
private final ParquetDataSource dataSource;
private final DateTimeZone timeZone;
private final AggregatedMemoryContext systemMemoryContext;
private int currentBlock;
@ -87,11 +89,13 @@ public class ParquetReader
public ParquetReader(MessageColumnIO messageColumnIO,
List<BlockMetaData> blocks,
ParquetDataSource dataSource,
DateTimeZone timeZone,
AggregatedMemoryContext systemMemoryContext,
DataSize maxReadBlockSize)
{
this.blocks = blocks;
this.dataSource = requireNonNull(dataSource, "dataSource is null");
this.timeZone = requireNonNull(timeZone, "timeZone is null");
this.systemMemoryContext = requireNonNull(systemMemoryContext, "systemMemoryContext is null");
this.currentRowGroupMemoryContext = systemMemoryContext.newAggregatedMemoryContext();
this.maxReadBlockBytes = requireNonNull(maxReadBlockSize, "maxReadBlockSize is null").toBytes();
@ -257,7 +261,7 @@ public class ParquetReader
{
for (PrimitiveColumnIO columnIO : columns) {
RichColumnDescriptor column = new RichColumnDescriptor(columnIO.getColumnDescriptor(), columnIO.getType().asPrimitiveType());
columnReaders[columnIO.getId()] = PrimitiveColumnReader.createReader(column);
columnReaders[columnIO.getId()] = PrimitiveColumnReader.createReader(column, timeZone);
}
}

View File

@ -35,6 +35,7 @@ import org.apache.parquet.column.ColumnDescriptor;
import org.apache.parquet.column.values.ValuesReader;
import org.apache.parquet.column.values.rle.RunLengthBitPackingHybridDecoder;
import org.apache.parquet.io.ParquetDecodingException;
import org.joda.time.DateTimeZone;
import java.io.ByteArrayInputStream;
import java.io.IOException;
@ -80,7 +81,7 @@ public abstract class PrimitiveColumnReader
return ParquetTypeUtils.isValueNull(columnDescriptor.isRequired(), definitionLevel, columnDescriptor.getMaxDefinitionLevel());
}
public static PrimitiveColumnReader createReader(RichColumnDescriptor descriptor)
public static PrimitiveColumnReader createReader(RichColumnDescriptor descriptor, DateTimeZone timeZone)
{
switch (descriptor.getType()) {
case BOOLEAN:
@ -90,7 +91,7 @@ public abstract class PrimitiveColumnReader
case INT64:
return createDecimalColumnReader(descriptor).orElse(new LongColumnReader(descriptor));
case INT96:
return new TimestampColumnReader(descriptor);
return new TimestampColumnReader(descriptor, timeZone);
case FLOAT:
return new FloatColumnReader(descriptor);
case DOUBLE:

View File

@ -15,25 +15,38 @@ package io.prestosql.parquet.reader;
import io.prestosql.parquet.RichColumnDescriptor;
import io.prestosql.spi.block.BlockBuilder;
import io.prestosql.spi.type.TimestampWithTimeZoneType;
import io.prestosql.spi.type.Type;
import org.apache.parquet.io.api.Binary;
import org.joda.time.DateTimeZone;
import static io.prestosql.parquet.ParquetTimestampUtils.getTimestampMillis;
import static io.prestosql.spi.type.DateTimeEncoding.packDateTimeWithZone;
import static io.prestosql.spi.type.TimeZoneKey.UTC_KEY;
import static java.util.Objects.requireNonNull;
public class TimestampColumnReader
extends PrimitiveColumnReader
{
public TimestampColumnReader(RichColumnDescriptor descriptor)
private final DateTimeZone timeZone;
public TimestampColumnReader(RichColumnDescriptor descriptor, DateTimeZone timeZone)
{
super(descriptor);
this.timeZone = requireNonNull(timeZone, "timeZone is null");
}
@Override
protected void readValue(BlockBuilder blockBuilder, Type type)
{
if (definitionLevel == columnDescriptor.getMaxDefinitionLevel()) {
Binary binary = valuesReader.readBytes();
type.writeLong(blockBuilder, getTimestampMillis(binary));
long utcMillis = getTimestampMillis(valuesReader.readBytes());
if (type instanceof TimestampWithTimeZoneType) {
type.writeLong(blockBuilder, packDateTimeWithZone(utcMillis, UTC_KEY));
}
else {
utcMillis = timeZone.convertUTCToLocal(utcMillis);
type.writeLong(blockBuilder, utcMillis);
}
}
else if (isValueNull()) {
blockBuilder.appendNull();