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
12 changes: 6 additions & 6 deletions amber/src/main/python/core/models/schema/attribute_type.py
Original file line number Diff line number Diff line change
Expand Up @@ -98,18 +98,18 @@ def _parse_bool(v):


def _parse_timestamp(v):
# A TIMESTAMP reaches a worker as a wall clock with no zone, so a value
# parsed here holds none either and compares with a row's. An offset is
# moved to the machine's zone first, as the engine's DateParserUtils does.
if _is_empty_value(v):
return datetime.datetime(1970, 1, 1, tzinfo=datetime.timezone.utc)
return datetime.datetime(1970, 1, 1)

normalized_value = str(v)
if normalized_value.endswith("Z"):
normalized_value = normalized_value[:-1] + "+00:00"
parsed_value = datetime.datetime.fromisoformat(normalized_value)
if (
parsed_value.tzinfo is None
or parsed_value.tzinfo.utcoffset(parsed_value) is None
):
return parsed_value.replace(tzinfo=datetime.timezone.utc)
if parsed_value.utcoffset() is not None:
return parsed_value.astimezone().replace(tzinfo=None)
return parsed_value


Expand Down
57 changes: 30 additions & 27 deletions amber/src/test/python/core/models/schema/test_attribute_type.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
# under the License.

import datetime
import time

import pytest

Expand All @@ -31,7 +32,19 @@
parse_bool = FROM_STRING_PARSER_MAPPING[AttributeType.BOOL]
parse_timestamp = FROM_STRING_PARSER_MAPPING[AttributeType.TIMESTAMP]

EPOCH = datetime.datetime(1970, 1, 1, tzinfo=datetime.timezone.utc)
EPOCH = datetime.datetime(1970, 1, 1)


@pytest.fixture
def machine_zone(monkeypatch):
# An offset is moved to the machine's zone, so the tests set one.
def set_zone(name):
monkeypatch.setenv("TZ", name)
time.tzset()

yield set_zone
monkeypatch.undo()
time.tzset()


class TestParseBool:
Expand Down Expand Up @@ -60,35 +73,25 @@ def test_a_non_numeric_non_literal_value_is_rejected(self):

class TestParseTimestamp:
@pytest.mark.parametrize("empty", [None, "", " ", "\t\n"])
def test_an_absent_value_parses_as_the_utc_epoch(self, empty):
def test_an_absent_value_parses_as_the_epoch_with_no_zone(self, empty):
parsed = parse_timestamp(empty)
# Assert the exact instant *and* the tzinfo: a naive 1970-01-01 would
# compare unequal here, so dropping the timezone is caught too.
assert parsed == EPOCH
assert parsed.tzinfo == datetime.timezone.utc
assert (parsed.year, parsed.month, parsed.day) == (1970, 1, 1)

def test_a_zulu_suffix_yields_a_utc_aware_instant(self):
# Named for the observable outcome, not for the code that produces it.
# `_parse_timestamp` rewrites a trailing "Z" into "+00:00" before
# calling `fromisoformat`, but `fromisoformat` has accepted "Z" itself
# since Python 3.11 and the pyamber CI matrix is 3.11/3.12/3.13 -- so
# deleting that rewrite leaves the entire suite green, and this test
# must not be credited with pinning it. What it does pin is the
# resulting instant and its offset (and, uniquely in this file, the
# "+00:00" constant, should the rewrite ever run on an older runtime).
assert parse_timestamp("2024-05-06T07:08:09Z") == datetime.datetime(
2024, 5, 6, 7, 8, 9, tzinfo=datetime.timezone.utc
)
assert parsed.tzinfo is None

def test_a_value_with_no_offset_keeps_its_wall_clock(self, machine_zone):
machine_zone("Asia/Kolkata")
parsed = parse_timestamp("2024-05-06T07:08:09")
assert parsed == datetime.datetime(2024, 5, 6, 7, 8, 9)
assert parsed.tzinfo is None

def test_a_naive_value_is_assumed_to_be_utc(self):
assert parse_timestamp("2024-05-06T07:08:09") == datetime.datetime(
2024, 5, 6, 7, 8, 9, tzinfo=datetime.timezone.utc
def test_a_zulu_suffix_is_moved_to_the_machine_zone(self, machine_zone):
machine_zone("Asia/Kolkata")
assert parse_timestamp("2024-05-06T07:08:09Z") == datetime.datetime(
2024, 5, 6, 12, 38, 9
)

def test_an_explicit_offset_is_preserved(self):
def test_an_explicit_offset_is_moved_to_the_machine_zone(self, machine_zone):
machine_zone("UTC")
parsed = parse_timestamp("2024-05-06T07:08:09+02:00")
assert parsed.utcoffset() == datetime.timedelta(hours=2)
assert parsed == datetime.datetime(
2024, 5, 6, 5, 8, 9, tzinfo=datetime.timezone.utc
)
assert parsed == datetime.datetime(2024, 5, 6, 5, 8, 9)
assert parsed.tzinfo is None
9 changes: 2 additions & 7 deletions amber/src/test/python/pytexera/udf/test_udf_operator.py
Original file line number Diff line number Diff line change
Expand Up @@ -206,7 +206,7 @@ def test_injected_values_are_applied_before_open(self):
assert operator.count_parameter.value == 7
assert operator.enabled_parameter.value is True
assert operator.created_at_parameter.value == datetime.datetime(
2024, 1, 1, 0, 0, tzinfo=datetime.timezone.utc
2024, 1, 1, 0, 0
)

def test_duplicate_parameter_names_with_conflicting_types_raise(self):
Expand Down Expand Up @@ -238,12 +238,7 @@ def test_duplicate_parameter_names_with_same_type_succeed(self):
(
"2024-01-01T00:00:00",
AttributeType.TIMESTAMP,
datetime.datetime(2024, 1, 1, 0, 0, tzinfo=datetime.timezone.utc),
),
(
"2024-01-01T00:00:00Z",
AttributeType.TIMESTAMP,
datetime.datetime(2024, 1, 1, 0, 0, tzinfo=datetime.timezone.utc),
datetime.datetime(2024, 1, 1, 0, 0),
),
],
)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -144,8 +144,14 @@ object AttributeTypeUtils extends Serializable {
} else {
str.trim.toInt
}
case int: Integer => int
case long: java.lang.Long => long.toInt
case int: Integer => int
case long: java.lang.Long => long.toInt
// An Arrow file states the width of its own integers, and a narrow
// column hands its values over as Shorts or Bytes. Without this the parse
// threw, the Arrow source caught it, and every value in the column
// arrived null.
case short: java.lang.Short => short.toInt
case byte: java.lang.Byte => byte.toInt
case double: java.lang.Double => double.toInt
case boolean: java.lang.Boolean => if (boolean) 1 else 0
// Timestamp and Binary are considered to be illegal here.
Expand Down Expand Up @@ -232,10 +238,15 @@ object AttributeTypeUtils extends Serializable {
def parseDouble(fieldValue: Any): java.lang.Double = {
val attempt: Try[Double] = Try {
fieldValue match {
case str: String => str.trim.toDouble
case int: Integer => int.toDouble
case long: java.lang.Long => long.toDouble
case double: java.lang.Double => double
case str: String => str.trim.toDouble
case int: Integer => int.toDouble
case long: java.lang.Long => long.toDouble
case double: java.lang.Double => double
// A single-precision column hands its values over as Floats, which only
// an Arrow file produces: Texera writes every double it owns as eight
// bytes. Without this the parse threw, the Arrow source caught it, and
// the whole column arrived null.
case float: java.lang.Float => float.toDouble
case boolean: java.lang.Boolean => if (boolean) 1 else 0
// Timestamp and Binary are considered to be illegal here.
case _ =>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -84,11 +84,17 @@ object ArrowUtils extends LazyLogging {
// Use the attribute type from the schema (which includes metadata)
// instead of deriving it from the Arrow type
val attributeType = schema.getAttributes(index).getType
// A timestamp is the one type whose field says more than the
// schema does, so it is the only one that reads the field.
// A timestamp and an unsigned integer are the types whose field
// says more than the schema does, so they are the only ones that
// read the field: an unsigned column arrives as INTEGER or LONG.
attributeType match {
case AttributeType.TIMESTAMP => wallClockOf(value, fieldVector.getField.getType)
case _ => AttributeTypeUtils.parseField(value, attributeType)
case AttributeType.INTEGER | AttributeType.LONG =>
AttributeTypeUtils.parseField(
unsignedValueOf(value, fieldVector.getField.getType),
attributeType
)
case _ => AttributeTypeUtils.parseField(value, attributeType)
}
} catch {
case e: Exception =>
Expand All @@ -101,6 +107,29 @@ object ArrowUtils extends LazyLogging {
.build()
}

/** The number an unsigned column counts to, out of the storage it counts in.
*
* Arrow's unsigned vectors hand back the raw storage, signed: `UInt1Vector` a
* Byte, `UInt4Vector` an Int, each reading the file's largest value as -1.
* `UInt2Vector` is the exception and hands back a Character, which is already
* the number but is no kind of integer to the parse below. Every other value
* passes through untouched.
*
* Nothing wider is handled because nothing wider arrives: [[toAttributeType]]
* refuses an unsigned 64-bit column, having no Texera type to hold it.
*/
private def unsignedValueOf(value: AnyRef, arrowType: ArrowType): AnyRef =
arrowType match {
case int: ArrowType.Int if !int.getIsSigned =>
value match {
case byte: java.lang.Byte => Int.box(byte & 0xff)
case char: java.lang.Character => Int.box(char.toInt)
case integer: Integer => Long.box(integer.toLong & 0xffffffffL)
case other => other
}
case _ => value
}

/** The wall clock a timestamp column holds, read the way its own field states.
*
* A Texera TIMESTAMP carries no zone, so the wall clock is the whole of what
Expand Down Expand Up @@ -140,13 +169,17 @@ object ArrowUtils extends LazyLogging {
case TimeUnit.NANOSECOND => Instant.EPOCH.plusNanos(number)
}

/** The number a field of this unit records an instant as, inverting [[instantOf]]. */
/** The number a field of this unit records an instant as, inverting [[instantOf]].
* Microseconds are counted from the seconds, because ChronoUnit.MICROS goes
* through nanoseconds and overflows past 2262, a year a Timestamp holds.
*/
private def numberOf(instant: Instant, unit: TimeUnit): Long =
unit match {
case TimeUnit.SECOND => instant.getEpochSecond
case TimeUnit.MILLISECOND => instant.toEpochMilli
case TimeUnit.MICROSECOND => ChronoUnit.MICROS.between(Instant.EPOCH, instant)
case TimeUnit.NANOSECOND => ChronoUnit.NANOS.between(Instant.EPOCH, instant)
case TimeUnit.MICROSECOND =>
Math.addExact(Math.multiplyExact(instant.getEpochSecond, 1000000L), instant.getNano / 1000L)
case TimeUnit.NANOSECOND => ChronoUnit.NANOS.between(Instant.EPOCH, instant)
}

/**
Expand Down Expand Up @@ -182,15 +215,23 @@ object ArrowUtils extends LazyLogging {
@throws[AttributeTypeException]
def toAttributeType(srcType: ArrowType): AttributeType = {
srcType match {
// An unsigned column counts up where its storage counts down: the largest
// unsigned 32-bit value is stored as -1, so read as its storage it would
// arrive as -1 where the file means 4294967295. The next Texera integer up
// holds it, and [[unsignedValueOf]] does the reading. Past 64 bits there is
// no next one. The same widening covers the narrow widths, Texera having no
// column shorter than a 32-bit integer.
case int: ArrowType.Int =>
int.getBitWidth match {
case 16 | 32 =>
AttributeType.INTEGER

case 64 =>
AttributeType.LONG

case other =>
(int.getBitWidth, int.getIsSigned) match {
case (8 | 16 | 32, true) => AttributeType.INTEGER
case (8 | 16, false) => AttributeType.INTEGER
case (64, true) => AttributeType.LONG
case (32, false) => AttributeType.LONG
case (64, false) =>
throw new AttributeTypeUtils.AttributeTypeException(
"Unsupported unsigned 64-bit Int, which is wider than any Texera column"
)
case (other, _) =>
throw new AttributeTypeUtils.AttributeTypeException(
s"Unsupported Int bit width: $other"
)
Expand Down Expand Up @@ -366,8 +407,13 @@ object ArrowUtils extends LazyLogging {
case AttributeType.BOOLEAN =>
ArrowType.Bool.INSTANCE

// A wall clock with no zone, to the microsecond, which is what a Python
// worker maps TIMESTAMP to and the finest a datetime holds. Labelled UTC
// in milliseconds, a Python operator was handed an aware moment cut to
// the millisecond, where a JVM operator reads the same row with no zone
// and to the nanosecond.
case AttributeType.TIMESTAMP =>
new ArrowType.Timestamp(TimeUnit.MILLISECOND, "UTC")
new ArrowType.Timestamp(TimeUnit.MICROSECOND, null)

case AttributeType.BINARY =>
new ArrowType.Binary
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -50,15 +50,27 @@ class ArrowUtilsSpec extends AnyFlatSpec with Matchers {
ArrowUtils.toAttributeType(new ArrowType.Int(64, true)) shouldBe AttributeType.LONG
}

it should "throw AttributeTypeException for non-standard Int bit-widths" in {
// Only 16/32 (INTEGER) and 64 (LONG) are supported. Other widths used to
// be silently coerced to LONG by a `case 64 | _` catch-all; they now
// raise rather than masquerade as Int64.
it should "map every integer to the narrowest Texera type that holds it" in {
// Texera has no column shorter than a 32-bit integer, so a narrower width
// reads as one. An unsigned column needs the width above its own, counting
// up where its storage counts down: the largest unsigned 32-bit value is
// past what an INTEGER holds.
ArrowUtils.toAttributeType(new ArrowType.Int(8, true)) shouldBe AttributeType.INTEGER
ArrowUtils.toAttributeType(new ArrowType.Int(8, false)) shouldBe AttributeType.INTEGER
ArrowUtils.toAttributeType(new ArrowType.Int(16, false)) shouldBe AttributeType.INTEGER
ArrowUtils.toAttributeType(new ArrowType.Int(32, false)) shouldBe AttributeType.LONG
}

it should "throw AttributeTypeException for an integer no Texera type holds" in {
// A width above 64 used to be silently coerced to LONG by a `case 64 | _`
// catch-all; it raises rather than masquerade as Int64. An unsigned 64-bit
// column raises for the same reason, there being no width above it to read
// it as.
assertThrows[AttributeTypeException] {
ArrowUtils.toAttributeType(new ArrowType.Int(8, true))
ArrowUtils.toAttributeType(new ArrowType.Int(128, true))
}
assertThrows[AttributeTypeException] {
ArrowUtils.toAttributeType(new ArrowType.Int(128, true))
ArrowUtils.toAttributeType(new ArrowType.Int(64, false))
}
}

Expand Down Expand Up @@ -114,9 +126,9 @@ class ArrowUtilsSpec extends AnyFlatSpec with Matchers {
ArrowUtils.fromAttributeType(AttributeType.BOOLEAN) shouldBe ArrowType.Bool.INSTANCE
}

it should "map TIMESTAMP to Timestamp(MILLISECOND, UTC)" in {
it should "map TIMESTAMP to a zoneless Timestamp(MICROSECOND)" in {
val arrow = ArrowUtils.fromAttributeType(AttributeType.TIMESTAMP)
arrow shouldBe new ArrowType.Timestamp(TimeUnit.MILLISECOND, "UTC")
arrow shouldBe new ArrowType.Timestamp(TimeUnit.MICROSECOND, null)
}

it should "map BINARY to ArrowType.Binary" in {
Expand Down Expand Up @@ -357,6 +369,20 @@ class ArrowUtilsSpec extends AnyFlatSpec with Matchers {
}
}

// A Timestamp reaches past 2262, where counting microseconds through
// nanoseconds overflowed, and before 1970. Both come back to the microsecond.
it should "round-trip a timestamp either side of the nanosecond range" in {
val schema = Schema(List(new Attribute("t", AttributeType.TIMESTAMP)))
Seq("2500-01-01 00:00:00.123456", "1600-06-15 08:30:00.5").foreach { text =>
val tuple =
Tuple.builder(schema).addSequentially(Array[Any](Timestamp.valueOf(text))).build()
withRoot(schema) { root =>
ArrowUtils.appendTexeraTuple(tuple, root)
ArrowUtils.getTexeraTuple(0, root).getField[Timestamp]("t") shouldBe Timestamp.valueOf(text)
}
}
}

it should "append consecutive tuples at increasing row indices" in {
val schema = Schema(List(new Attribute("s", AttributeType.STRING)))
val first = Tuple.builder(schema).addSequentially(Array[Any]("first")).build()
Expand Down Expand Up @@ -436,9 +462,9 @@ class ArrowUtilsSpec extends AnyFlatSpec with Matchers {
}
}

// ----- Timestamp fields that are not the UTC millisecond ones we write -----
// ----- Timestamp fields that are not the zoneless microsecond ones we write -----

// fromTexeraSchema only ever writes Timestamp(MILLISECOND, "UTC"), so the
// fromTexeraSchema only ever writes a zoneless Timestamp(MICROSECOND), so the
// roots built above never exercise another zone or unit. An .arrow file handed
// to ArrowSourceOpDesc can carry either: pandas writes a tz-aware column as
// Timestamp(NANOSECOND, <its zone>). These build the field directly.
Expand Down
Loading
Loading