Skip to content
Closed
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 @@ -25,11 +25,14 @@
import java.time.Duration;
import java.time.Instant;
import java.time.LocalDate;
import java.time.LocalTime;
import java.time.OffsetTime;
import java.time.Period;
import java.time.ZoneId;
import java.time.ZoneOffset;
import java.time.ZonedDateTime;
import java.time.format.DateTimeFormatter;
import java.time.temporal.ChronoUnit;
import java.util.HashMap;
import java.util.Iterator;
import java.util.List;
Expand Down Expand Up @@ -563,6 +566,10 @@ static void handleLogicalFieldType(
(ByteBuffer) element.get(fieldName), fieldSchema, fieldSchema.getLogicalType());
jsonObject.put(fieldName, bigDecimal.toPlainString());
} else if (fieldSchema.getLogicalType() instanceof LogicalTypes.TimeMicros) {
/*
* Note: Upstream Datastream incorrectly wraps 24:00:00 to 0 microseconds during extraction.
* This causes 24:00:00 to be extracted as 00:00:00 (e.g. 'PT0S' instead of 'PT24H').
*/
Long microseconds = (Long) element.get(fieldName);
if (microseconds.equals(DATETIME_POSITIVE_INFINITY)) {
jsonObject.put(fieldName, "infinity");
Expand Down Expand Up @@ -650,6 +657,58 @@ static void handleDatastreamRecordType(
.withZoneSameInstant(ZoneId.of("UTC"))
.format(DEFAULT_TIMESTAMP_WITH_TZ_FORMATTER));
break;
case "timeTz":
/*
* Note: Upstream Datastream incorrectly wraps 24:00:00 to 0 microseconds during extraction.
* This causes 24:00:00+offset to be extracted as 00:00:00+offset.
*/
long timeTzNanos =
((Number) getOrDefault(element, "time", 0L)).longValue()
* TimeUnit.MICROSECONDS.toNanos(1);
int offsetSeconds = ((Number) getOrDefault(element, "offset", 0)).intValue() / 1000;

ZoneOffset timeTzOffset = ZoneOffset.ofTotalSeconds(offsetSeconds);

if (timeTzNanos == 86400000000000L) {
jsonObject.put(fieldName, "24:00:00" + timeTzOffset.toString());
break;
}

LocalTime localTime = LocalTime.ofNanoOfDay(timeTzNanos);
OffsetTime offsetTime = OffsetTime.of(localTime, timeTzOffset);
jsonObject.put(fieldName, offsetTime.format(DateTimeFormatter.ISO_OFFSET_TIME));
break;
case "interval":
int months = ((Number) getOrDefault(element, "months", 0)).intValue();
int hours = ((Number) getOrDefault(element, "hours", 0)).intValue();
long micros = ((Number) getOrDefault(element, "micros", 0L)).longValue();

int days = hours / 24;
int remainingHours = hours % 24;

Period intervalPeriod = Period.ZERO.plusMonths(months).plusDays(days).normalized();
Duration intervalDuration =
Duration.ZERO.plusHours(remainingHours).plus(micros, ChronoUnit.MICROS);

if (intervalPeriod.isZero() && intervalDuration.isZero()) {
jsonObject.put(fieldName, "PT0S");
break;
}

StringBuilder result = new StringBuilder();
if (!intervalPeriod.isZero()) {
result.append(intervalPeriod.toString());
}
if (!intervalDuration.isZero()) {
if (result.length() == 0) {
result.append(intervalDuration.toString());
} else {
// Remove the 'P' from Duration
result.append(intervalDuration.toString().substring(1));
}
}
jsonObject.put(fieldName, result.toString());
break;
/*
* The `intervalNano` maps to nano second precision interval type used by Cassandra Interval.
* On spanner this will map to `string` or `Interval` type.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -197,7 +197,7 @@ public void testPostgresByteArray() throws IOException, URISyntaxException {
ObjectMapper mapper = new ObjectMapper();
JsonNode changeEvent = mapper.readTree(jsonData);
// The avro file contains binary_content: b'\xde\xad\xbe\xef', which is converted to
// base64 encoded string by Jackson library.
// base64 encoded string by Jackson.
assertEquals("3q2+7w==", changeEvent.get("binary_content").textValue());
}

Expand Down Expand Up @@ -316,6 +316,76 @@ public void testIntervalNano() throws JsonProcessingException {
assertEquals(expected, new ObjectMapper().writeValueAsString(objectNode));
}

@Test
public void testInterval() throws JsonProcessingException {
ObjectNode objectNode = new ObjectNode(new JsonNodeFactory(true));

/* Basic Test: 1 month, 26 hours (1 day + 2 hours), 3000000 micros (3 seconds) */
UnifiedTypesFormatter.handleDatastreamRecordType(
"basic", generateIntervalSchema(), generateIntervalRecord(1, 26, 3000000L), objectNode);

/* Zero interval */
UnifiedTypesFormatter.handleDatastreamRecordType(
"zero_interval", generateIntervalSchema(), generateIntervalRecord(0, 0, 0L), objectNode);

/* Only months */
UnifiedTypesFormatter.handleDatastreamRecordType(
"only_months", generateIntervalSchema(), generateIntervalRecord(12, 0, 0L), objectNode);

/* Only time (5 hours + 123456 micros = 5H0.123456S) */
UnifiedTypesFormatter.handleDatastreamRecordType(
"only_time", generateIntervalSchema(), generateIntervalRecord(0, 5, 123456L), objectNode);

/* Negative interval (-1 month, -26 hours = -1 day - 2 hours, -3000000 micros = -3 seconds) */
UnifiedTypesFormatter.handleDatastreamRecordType(
"neg_basic",
generateIntervalSchema(),
generateIntervalRecord(-1, -26, -3000000L),
objectNode);

String expected =
"{\"basic\":\"P1M1DT2H3S\","
+ "\"zero_interval\":\"PT0S\","
+ "\"only_months\":\"P1Y\","
+ "\"only_time\":\"PT5H0.123456S\","
+ "\"neg_basic\":\"P-1M-1DT-2H-3S\"}";
assertEquals(expected, new ObjectMapper().writeValueAsString(objectNode));
}

@Test
public void testTimeTz() throws JsonProcessingException {
ObjectNode objectNode = new ObjectNode(new JsonNodeFactory(true));

/* Basic Test: 23:59:59 + 10 hours offset */
UnifiedTypesFormatter.handleDatastreamRecordType(
"basic", generateTimeTzSchema(), generateTimeTzRecord(86399000000L, 36000000), objectNode);

/* Negative offset: 12:30:00 - 5 hours offset */
UnifiedTypesFormatter.handleDatastreamRecordType(
"neg_offset",
generateTimeTzSchema(),
generateTimeTzRecord(45000000000L, -18000000),
objectNode);

/* Zero offset (UTC): 08:00:00Z */
UnifiedTypesFormatter.handleDatastreamRecordType(
"utc", generateTimeTzSchema(), generateTimeTzRecord(28800000000L, 0), objectNode);

/* 24:00:00 special case */
UnifiedTypesFormatter.handleDatastreamRecordType(
"max_time",
generateTimeTzSchema(),
generateTimeTzRecord(86400000000L, 36000000),
objectNode);

String expected =
"{\"basic\":\"23:59:59+10:00\","
+ "\"neg_offset\":\"12:30:00-05:00\","
+ "\"utc\":\"08:00:00Z\","
+ "\"max_time\":\"24:00:00+10:00\"}";
assertEquals(expected, new ObjectMapper().writeValueAsString(objectNode));
}

@Test
public void testGetPrimaryKeys_primaryKeysField() throws IOException {
Schema arraySchema = Schema.createArray(Schema.create(Schema.Type.STRING));
Expand Down Expand Up @@ -570,4 +640,48 @@ private GenericRecord buildOuterRecord(GenericRecord sourceMetadata, String read
record.put("payload", payload);
return record;
}

private GenericRecord generateIntervalRecord(Integer months, Integer hours, Long micros) {
GenericRecord genericRecord = new GenericData.Record(generateIntervalSchema());
genericRecord.put("months", months);
genericRecord.put("hours", hours);
genericRecord.put("micros", micros);
return genericRecord;
}

private Schema generateIntervalSchema() {
return SchemaBuilder.builder()
.record("interval")
.fields()
.name("months")
.type(SchemaBuilder.builder().intType())
.withDefault(0)
.name("hours")
.type(SchemaBuilder.builder().intType())
.withDefault(0)
.name("micros")
.type(SchemaBuilder.builder().longType())
.withDefault(0L)
.endRecord();
}

private GenericRecord generateTimeTzRecord(Long time, Integer offset) {
GenericRecord genericRecord = new GenericData.Record(generateTimeTzSchema());
genericRecord.put("time", time);
genericRecord.put("offset", offset);
return genericRecord;
}

private Schema generateTimeTzSchema() {
return SchemaBuilder.builder()
.record("timeTz")
.fields()
.name("time")
.type(SchemaBuilder.builder().longType())
.withDefault(0L)
.name("offset")
.type(SchemaBuilder.builder().intType())
.withDefault(0)
.endRecord();
}
}
Loading
Loading