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
75 changes: 75 additions & 0 deletions beam/src/main/scala/magnolify/beam/logical/CompatTypes.scala
Original file line number Diff line number Diff line change
@@ -0,0 +1,75 @@
/*
* Copyright 2026 Spotify AB
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/

package magnolify.beam.logical

import magnolify.beam.RowField
import magnolify.shared.Time._
import org.apache.beam.sdk.schemas.Schema.FieldType
import org.apache.beam.sdk.schemas.logicaltypes
import org.joda.time as joda
import org.joda.time.chrono.ISOChronology

import java.time as jt

// The pre-0.9.8 `Instant` encodings, shared by `compat.*` and by the deprecated bare
// `millis`/`micros`/`nanos` objects so that the two cannot drift apart.
//
// Note these three do not share a representation -- millis is the joda-backed `DATETIME`
// primitive, micros is a raw `INT64` and nanos is the SDK-local `NanosInstant` logical type.
// What they have in common is only that this is what 0.9.7 produced, which is why the grouping
// is named for its history rather than for an encoding.

private[logical] trait MillisCompat extends MillisNonInstant {
// Layered on the joda mapping below, so the typeclass's internal representation is joda.
// `lazy` because it resolves `rfJodaInstantMillis`, which is declared after it.
implicit lazy val rfInstantMillis: RowField[jt.Instant] =
RowField.from[joda.Instant](i => millisToInstant(millisFromJodaInstant(i)))(i =>
millisToJodaInstant(millisFromInstant(i))
)
implicit val rfJodaInstantMillis: RowField[joda.Instant] =
RowField.id[joda.Instant](_ => FieldType.DATETIME)
implicit val rfJodaDateTimeMillis: RowField[joda.DateTime] =
RowField.from[joda.Instant](_.toDateTime(ISOChronology.getInstanceUTC))(_.toInstant)
}

private[logical] trait MicrosCompat extends MicrosNonInstant {
// NOTE: logicaltypes.MicrosInstant() cannot be used as it throws assertion
// errors when greater-than-microsecond precision data is used
implicit val rfInstantMicros: RowField[jt.Instant] =
RowField.from[Long](microsToInstant)(microsFromInstant)
// joda.Instant has millisecond precision, excess precision discarded
implicit val rfJodaInstantMicros: RowField[joda.Instant] =
RowField.from[Long](microsToJodaInstant)(microsFromJodaInstant)
// joda.DateTime only has millisecond resolution, so excess precision is discarded
implicit val rfJodaDateTimeMicros: RowField[joda.DateTime] =
RowField.from[Long](microsToJodaDateTime)(microsFromJodaDateTime)
}

private[logical] trait NanosCompat extends NanosNonInstant {
implicit val rfInstantNanos: RowField[jt.Instant] =
RowField.id[jt.Instant](_ => FieldType.logicalType(new logicaltypes.NanosInstant()))
// joda.Instant has millisecond precision, excess precision discarded
implicit val rfJodaInstantNanos: RowField[joda.Instant] =
RowField.from[jt.Instant](i => nanosToJodaInstant(nanosFromInstant(i)))(i =>
nanosToInstant(nanosFromJodaInstant(i))
)
// joda.DateTime only has millisecond resolution
implicit val rfJodaDateTimeNanos: RowField[joda.DateTime] =
RowField.from[jt.Instant](i => nanosToJodaDateTime(nanosFromInstant(i)))(i =>
nanosToInstant(nanosFromJodaDateTime(i))
)
}
96 changes: 96 additions & 0 deletions beam/src/main/scala/magnolify/beam/logical/NonInstantTypes.scala
Original file line number Diff line number Diff line change
@@ -0,0 +1,96 @@
/*
* Copyright 2024 Spotify AB
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/

package magnolify.beam.logical

import magnolify.beam.RowField
import magnolify.shared.Time._
import org.apache.beam.sdk.schemas.Schema.FieldType
import org.apache.beam.sdk.schemas.logicaltypes
import org.joda.time as joda

import java.time as jt

// Mappings that Beam represents identically regardless of the instant encoding, shared by
// `timestamp.*`, `compat.*` and the deprecated bare `millis`/`micros`/`nanos` objects.
// "NonInstant" means invariant across those groupings at a given precision -- all six of
// these mappings do still vary by precision.

private[logical] trait MillisNonInstant {
implicit val rfLocalTimeMillis: RowField[jt.LocalTime] =
RowField.from[Int](millisToLocalTime)(millisFromLocalTime)
implicit val rfJodaLocalTimeMillis: RowField[joda.LocalTime] =
RowField.from[Int](millisToJodaLocalTime)(millisFromJodaLocalTime)

implicit val rfLocalDateTimeMillis: RowField[jt.LocalDateTime] =
RowField.id[jt.LocalDateTime](_ => FieldType.logicalType(new logicaltypes.DateTime()))
implicit val rfJodaLocalDateTimeMillis: RowField[joda.LocalDateTime] =
RowField.from[jt.LocalDateTime](ldt => millisToJodaLocalDateTime(millisFromLocalDateTime(ldt)))(
ldt => millisToLocalDateTime(millisFromJodaLocalDateTime(ldt))
)

implicit val rfDurationMillis: RowField[jt.Duration] =
RowField.from[Long](millisToDuration)(millisFromDuration)
implicit val rfJodaDurationMillis: RowField[joda.Duration] =
RowField.from[Long](millisToJodaDuration)(millisFromJodaDuration)
}

private[logical] trait MicrosNonInstant {
implicit val rfLocalTimeMicros: RowField[jt.LocalTime] =
RowField.from[Long](microsToLocalTime)(microsFromLocalTime)
// joda.LocalTime only has millisecond resolution, so excess precision is discarded
implicit val rfJodaLocalTimeMicros: RowField[joda.LocalTime] =
RowField.from[Long](microsToJodaLocalTime)(microsFromJodaLocalTime)

implicit val rfLocalDateTimeMicros: RowField[jt.LocalDateTime] =
RowField.from[Long](microsToLocalDateTime)(microsFromLocalDateTime)
// joda.LocalDateTime has millisecond precision, excess precision discarded
implicit val rfJodaLocalDateTimeMicros: RowField[joda.LocalDateTime] =
RowField.from[Long](microsToJodaLocalDateTime)(microsFromJodaLocalDateTime)

implicit val rfDurationMicros: RowField[jt.Duration] =
RowField.from[Long](microsToDuration)(microsFromDuration)
// joda.Duration has millisecond precision, excess precision discarded
implicit val rfJodaDurationMicros: RowField[joda.Duration] =
RowField.from[Long](microsToJodaDuration)(microsFromJodaDuration)
}

private[logical] trait NanosNonInstant {
implicit val rfLocalTimeNanos: RowField[jt.LocalTime] =
RowField.id[jt.LocalTime](_ => FieldType.logicalType(new logicaltypes.Time()))
// joda.LocalTime only has millisecond resolution, so excess precision is discarded
implicit val rfJodaLocalTimeNanos: RowField[joda.LocalTime] =
RowField.from[jt.LocalTime](lt => nanosToJodaLocalTime(nanosFromLocalTime(lt)))(lt =>
nanosToLocalTime(nanosFromJodaLocalTime(lt))
)

implicit val rfLocalDateTimeNanos: RowField[jt.LocalDateTime] =
RowField.from[Long](nanosToLocalDateTime)(nanosFromLocalDateTime)
// joda.LocalDateTime has millisecond precision, excess precision discarded
// NOTE: misnamed `Micros` since 0.9; kept for source/binary compatibility
implicit val rfJodaLocalDateTimeMicros: RowField[joda.LocalDateTime] =
RowField.from[jt.LocalDateTime](ldt => nanosToJodaLocalDateTime(nanosFromLocalDateTime(ldt)))(
ldt => nanosToLocalDateTime(nanosFromJodaLocalDateTime(ldt))
)

implicit val rfDurationNanos: RowField[jt.Duration] =
RowField.id[jt.Duration](_ => FieldType.logicalType(new logicaltypes.NanosDuration()))
// joda.Duration has millisecond precision, excess precision discarded
implicit val rfJodaDurationNanos: RowField[joda.Duration] =
RowField.from[jt.Duration](d => nanosToJodaDuration(nanosFromDuration(d)))(d =>
nanosToDuration(nanosFromJodaDuration(d))
)
}
Loading
Loading