Skip to content
Merged
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 @@ -650,6 +653,54 @@ static void handleDatastreamRecordType(
.withZoneSameInstant(ZoneId.of("UTC"))
.format(DEFAULT_TIMESTAMP_WITH_TZ_FORMATTER));
break;
case "timeTz":
Comment thread
shreyakhajanchi marked this conversation as resolved.
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":
Comment thread
shreyakhajanchi marked this conversation as resolved.
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;
Comment thread
jsuhani-2026 marked this conversation as resolved.
/*
* 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);

Comment thread
jsuhani-2026 marked this conversation as resolved.
/* 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