Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
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