Skip to content
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@
import java.time.Instant;
import java.time.LocalDate;
import java.time.LocalDateTime;
import java.time.LocalTime;
import java.time.ZoneId;
import java.time.ZoneOffset;
import java.time.temporal.ChronoUnit;
Expand All @@ -42,6 +43,7 @@
import static org.apache.flink.types.variant.BinaryVariantUtil.SIZE_LIMIT;
import static org.apache.flink.types.variant.BinaryVariantUtil.TIMESTAMP_FORMATTER;
import static org.apache.flink.types.variant.BinaryVariantUtil.TIMESTAMP_LTZ_FORMATTER;
import static org.apache.flink.types.variant.BinaryVariantUtil.TIME_FORMATTER;
import static org.apache.flink.types.variant.BinaryVariantUtil.VERSION;
import static org.apache.flink.types.variant.BinaryVariantUtil.VERSION_MASK;
import static org.apache.flink.types.variant.BinaryVariantUtil.checkIndex;
Expand Down Expand Up @@ -193,6 +195,26 @@ public Instant getInstant() throws VariantTypeException {
return microsToInstant(BinaryVariantUtil.getLong(value, pos));
}

@Override
public LocalTime getTime() throws VariantTypeException {
checkType(Type.TIME, getType());
return LocalTime.ofNanoOfDay(BinaryVariantUtil.getLong(value, pos) * 1000);
}

@Override
public LocalDateTime getDateTimeNanos() throws VariantTypeException {
checkType(Type.TIMESTAMP_NS, getType());
return nanosToInstant(BinaryVariantUtil.getLong(value, pos))
.atZone(ZoneOffset.UTC)
.toLocalDateTime();
}

@Override
public Instant getInstantNanos() throws VariantTypeException {
checkType(Type.TIMESTAMP_LTZ_NS, getType());
return nanosToInstant(BinaryVariantUtil.getLong(value, pos));
}

@Override
public byte[] getBytes() throws VariantTypeException {
checkType(Type.BYTES, getType());
Expand Down Expand Up @@ -224,10 +246,16 @@ public Object get() throws VariantTypeException {
return getString();
case DATE:
return getDate();
case TIME:
return getTime();
case TIMESTAMP:
return getDateTime();
case TIMESTAMP_LTZ:
return getInstant();
case TIMESTAMP_NS:
return getDateTimeNanos();
case TIMESTAMP_LTZ_NS:
return getInstantNanos();
case BYTES:
return getBytes();
default:
Expand Down Expand Up @@ -391,6 +419,27 @@ private static void toJsonImpl(
microsToInstant(BinaryVariantUtil.getLong(value, pos))
.atZone(ZoneOffset.UTC)));
break;
case TIME:
appendQuoted(
sb,
TIME_FORMATTER.format(
LocalTime.ofNanoOfDay(
BinaryVariantUtil.getLong(value, pos) * 1000)));
break;
case TIMESTAMP_LTZ_NS:
appendQuoted(
sb,
TIMESTAMP_LTZ_FORMATTER.format(
nanosToInstant(BinaryVariantUtil.getLong(value, pos))
.atZone(zoneId)));
break;
case TIMESTAMP_NS:
appendQuoted(
sb,
TIMESTAMP_FORMATTER.format(
nanosToInstant(BinaryVariantUtil.getLong(value, pos))
.atZone(ZoneOffset.UTC)));
break;
case FLOAT:
{
final float f = BinaryVariantUtil.getFloat(value, pos);
Expand Down Expand Up @@ -418,6 +467,10 @@ private static Instant microsToInstant(long timestamp) {
return Instant.EPOCH.plus(timestamp, ChronoUnit.MICROS);
}

private static Instant nanosToInstant(long timestamp) {
return Instant.EPOCH.plus(timestamp, ChronoUnit.NANOS);
}

private void checkType(Type expected, Type actual) {
if (expected != actual) {
throw new VariantTypeException(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@
import java.time.Instant;
import java.time.LocalDate;
import java.time.LocalDateTime;
import java.time.LocalTime;
import java.time.ZoneOffset;
import java.time.temporal.ChronoUnit;
import java.util.ArrayList;
Expand Down Expand Up @@ -106,7 +107,11 @@ public Variant of(BigDecimal bigDecimal) {
@Override
public Variant of(Instant instant) {
BinaryVariantInternalBuilder builder = new BinaryVariantInternalBuilder(false);
builder.appendTimestampLtz(ChronoUnit.MICROS.between(Instant.EPOCH, instant));
if (instant.getNano() % 1000 == 0) {
builder.appendTimestampLtz(ChronoUnit.MICROS.between(Instant.EPOCH, instant));
} else {
builder.appendTimestampLtzNanos(nanosSinceEpoch(instant));
}
return builder.build();
}

Expand All @@ -120,8 +125,32 @@ public Variant of(LocalDate localDate) {
@Override
public Variant of(LocalDateTime localDateTime) {
BinaryVariantInternalBuilder builder = new BinaryVariantInternalBuilder(false);
builder.appendTimestamp(
ChronoUnit.MICROS.between(Instant.EPOCH, localDateTime.toInstant(ZoneOffset.UTC)));
Instant instant = localDateTime.toInstant(ZoneOffset.UTC);
if (localDateTime.getNano() % 1000 == 0) {
builder.appendTimestamp(ChronoUnit.MICROS.between(Instant.EPOCH, instant));
} else {
builder.appendTimestampNanos(nanosSinceEpoch(instant));
}
return builder.build();
}

private static long nanosSinceEpoch(Instant instant) {
try {
return ChronoUnit.NANOS.between(Instant.EPOCH, instant);
} catch (ArithmeticException e) {
throw new VariantTypeException(
String.format(
"%s is outside the +/-292 year range (1677-09-21 to 2262-04-11) "
+ "supported by nanosecond precision variant timestamps. Use "
+ "microsecond precision instead.",
instant));
}
}

@Override
public Variant of(LocalTime localTime) {
BinaryVariantInternalBuilder builder = new BinaryVariantInternalBuilder(false);
builder.appendTime(localTime.toNanoOfDay() / 1000);
return builder.build();
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -58,8 +58,11 @@
import static org.apache.flink.types.variant.BinaryVariantUtil.NULL;
import static org.apache.flink.types.variant.BinaryVariantUtil.OBJECT;
import static org.apache.flink.types.variant.BinaryVariantUtil.SIZE_LIMIT;
import static org.apache.flink.types.variant.BinaryVariantUtil.TIME;
import static org.apache.flink.types.variant.BinaryVariantUtil.TIMESTAMP;
import static org.apache.flink.types.variant.BinaryVariantUtil.TIMESTAMP_LTZ;
import static org.apache.flink.types.variant.BinaryVariantUtil.TIMESTAMP_LTZ_NS;
import static org.apache.flink.types.variant.BinaryVariantUtil.TIMESTAMP_NS;
import static org.apache.flink.types.variant.BinaryVariantUtil.TRUE;
import static org.apache.flink.types.variant.BinaryVariantUtil.U16_MAX;
import static org.apache.flink.types.variant.BinaryVariantUtil.U24_MAX;
Expand Down Expand Up @@ -288,6 +291,27 @@ public void appendTimestamp(long microsSinceEpoch) {
writePos += 8;
}

public void appendTime(long microsSinceMidnight) {
checkCapacity(1 + 8);
writeBuffer[writePos++] = primitiveHeader(TIME);
writeLong(writeBuffer, writePos, microsSinceMidnight, 8);
writePos += 8;
}

public void appendTimestampLtzNanos(long nanosSinceEpoch) {
checkCapacity(1 + 8);
writeBuffer[writePos++] = primitiveHeader(TIMESTAMP_LTZ_NS);
writeLong(writeBuffer, writePos, nanosSinceEpoch, 8);
writePos += 8;
}

public void appendTimestampNanos(long nanosSinceEpoch) {
checkCapacity(1 + 8);
writeBuffer[writePos++] = primitiveHeader(TIMESTAMP_NS);
writeLong(writeBuffer, writePos, nanosSinceEpoch, 8);
writePos += 8;
}

public void appendFloat(float f) {
checkCapacity(1 + 4);
writeBuffer[writePos++] = primitiveHeader(FLOAT);
Expand Down
Loading