Timm0 commented on code in PR #29267: URL: https://github.com/apache/flink/pull/29267#discussion_r4143469847
########## flink-python/pyflink/table/typehints.py: ########## @@ -0,0 +1,160 @@ +################################################################################ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you 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. +################################################################################ +"""Shared inference from Python type hints to table :class:`DataType`s. + +This is the neutral core used to resolve a standard Python annotation into a +:mod:`pyflink.table.types` ``DataType``. It is consumed by the DataFrame API and +by UDF type-hint inference so both derive types from a single mapping. +""" + +import collections.abc +import datetime +import decimal +import types +from functools import partial +from typing import Any, Callable, Dict, Union, get_args, get_origin, get_type_hints + +from pyflink.table.types import DataType, DataTypes + +_PEP_604_UNION_TYPE = getattr(types, "UnionType", None) +_AWAITABLE_ORIGINS = (collections.abc.Coroutine, collections.abc.Awaitable) + +_BASIC_TYPE_HINT_FACTORIES: Dict[Any, Callable[[], DataType]] = { + bool: DataTypes.BOOLEAN, + int: DataTypes.BIGINT, + float: DataTypes.DOUBLE, + str: DataTypes.STRING, + bytes: DataTypes.BYTES, + bytearray: DataTypes.BYTES, + decimal.Decimal: partial(DataTypes.DECIMAL, 38, 18), + datetime.date: DataTypes.DATE, + # TIME is stored at runtime as an int number of milliseconds of the day, so 3 is + # the highest fractional-second precision that survives a round trip. + datetime.time: partial(DataTypes.TIME, 3), + datetime.datetime: DataTypes.TIMESTAMP, +} + + +def _is_typed_dict(type_hint: Any) -> bool: + try: + from typing import is_typeddict + + if is_typeddict(type_hint): + return True + except ImportError: + pass + return ( + isinstance(type_hint, type) + and issubclass(type_hint, dict) + and hasattr(type_hint, "__required_keys__") + ) + + +def _from_python_type(type_hint: Any) -> DataType: + """Resolve a Python type hint into a table :class:`DataType`. + + Supports the basic scalar types, ``list[T]``/``dict[K, V]`` containers, and + ``TypedDict`` (mapped to ``ROW``). Inferred types are ``NOT NULL``; + ``Optional[T]``/``T | None`` is the marker that widens a type to nullable, + and ``Any`` (an opt-out of the type system) stays nullable ``STRING``. + Raises :class:`TypeError` for hints that cannot be resolved unambiguously. + """ + + def infer_typed_dict(hint: Any) -> DataType: + return DataTypes.ROW( + [ + DataTypes.FIELD(name, infer(field_hint)) + for name, field_hint in get_type_hints(hint).items() Review Comment: Does this mean we silently ignore `typing.NotRequired` for `TypedDicts`? ########## flink-python/pyflink/table/typehints.py: ########## Review Comment: Could we add some tests for this new file? ########## flink-python/pyflink/table/typehints.py: ########## @@ -0,0 +1,160 @@ +################################################################################ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you 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. +################################################################################ +"""Shared inference from Python type hints to table :class:`DataType`s. + +This is the neutral core used to resolve a standard Python annotation into a +:mod:`pyflink.table.types` ``DataType``. It is consumed by the DataFrame API and +by UDF type-hint inference so both derive types from a single mapping. +""" + +import collections.abc +import datetime +import decimal +import types +from functools import partial +from typing import Any, Callable, Dict, Union, get_args, get_origin, get_type_hints + +from pyflink.table.types import DataType, DataTypes + +_PEP_604_UNION_TYPE = getattr(types, "UnionType", None) +_AWAITABLE_ORIGINS = (collections.abc.Coroutine, collections.abc.Awaitable) + +_BASIC_TYPE_HINT_FACTORIES: Dict[Any, Callable[[], DataType]] = { + bool: DataTypes.BOOLEAN, + int: DataTypes.BIGINT, + float: DataTypes.DOUBLE, + str: DataTypes.STRING, + bytes: DataTypes.BYTES, + bytearray: DataTypes.BYTES, + decimal.Decimal: partial(DataTypes.DECIMAL, 38, 18), + datetime.date: DataTypes.DATE, + # TIME is stored at runtime as an int number of milliseconds of the day, so 3 is + # the highest fractional-second precision that survives a round trip. + datetime.time: partial(DataTypes.TIME, 3), Review Comment: Setting a precision here makes sense, I just wonder if we should rather just mirror Java here where `LocalTime` is `TIME(9)` even though the highest fractional second precision is 3. ########## flink-python/pyflink/table/typehints.py: ########## @@ -0,0 +1,160 @@ +################################################################################ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you 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. +################################################################################ +"""Shared inference from Python type hints to table :class:`DataType`s. + +This is the neutral core used to resolve a standard Python annotation into a +:mod:`pyflink.table.types` ``DataType``. It is consumed by the DataFrame API and +by UDF type-hint inference so both derive types from a single mapping. +""" + +import collections.abc +import datetime +import decimal +import types +from functools import partial +from typing import Any, Callable, Dict, Union, get_args, get_origin, get_type_hints + +from pyflink.table.types import DataType, DataTypes + +_PEP_604_UNION_TYPE = getattr(types, "UnionType", None) +_AWAITABLE_ORIGINS = (collections.abc.Coroutine, collections.abc.Awaitable) + +_BASIC_TYPE_HINT_FACTORIES: Dict[Any, Callable[[], DataType]] = { + bool: DataTypes.BOOLEAN, + int: DataTypes.BIGINT, + float: DataTypes.DOUBLE, + str: DataTypes.STRING, + bytes: DataTypes.BYTES, + bytearray: DataTypes.BYTES, + decimal.Decimal: partial(DataTypes.DECIMAL, 38, 18), + datetime.date: DataTypes.DATE, + # TIME is stored at runtime as an int number of milliseconds of the day, so 3 is + # the highest fractional-second precision that survives a round trip. + datetime.time: partial(DataTypes.TIME, 3), + datetime.datetime: DataTypes.TIMESTAMP, +} + + +def _is_typed_dict(type_hint: Any) -> bool: + try: + from typing import is_typeddict + + if is_typeddict(type_hint): + return True + except ImportError: + pass + return ( + isinstance(type_hint, type) + and issubclass(type_hint, dict) + and hasattr(type_hint, "__required_keys__") + ) + + +def _from_python_type(type_hint: Any) -> DataType: + """Resolve a Python type hint into a table :class:`DataType`. + + Supports the basic scalar types, ``list[T]``/``dict[K, V]`` containers, and + ``TypedDict`` (mapped to ``ROW``). Inferred types are ``NOT NULL``; + ``Optional[T]``/``T | None`` is the marker that widens a type to nullable, + and ``Any`` (an opt-out of the type system) stays nullable ``STRING``. + Raises :class:`TypeError` for hints that cannot be resolved unambiguously. + """ + + def infer_typed_dict(hint: Any) -> DataType: + return DataTypes.ROW( Review Comment: Should we explicitly catch here a self-referring `TypedDict`? Currently that should just throw an error with an unintuitive error message. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
