MattBelle commented on code in PR #28797:
URL: https://github.com/apache/flink/pull/28797#discussion_r3633011430


##########
flink-python/pyflink/dataframe/convert.py:
##########
@@ -0,0 +1,232 @@
+################################################################################
+#  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.
+################################################################################
+
+from typing import Any, Iterable, Iterator, List, Mapping, Optional, Sequence, 
Tuple, Union, cast
+
+from pyflink.dataframe.context import get_or_create_table_environment
+from pyflink.dataframe.dataframe import DataFrame
+from pyflink.util.api_stability_decorators import PublicEvolving
+
+__all__ = ["from_dict", "from_records"]
+
+
+class _SequenceRowIterable:
+    """Produce tuple rows from row-oriented sequence records."""
+
+    def __init__(self, records: Sequence[Sequence[Any]]):
+        self._records = records
+
+    def __iter__(self) -> Iterator[Sequence[Any]]:
+        return (tuple(record) for record in self._records)
+
+
+class _ColumnRowIterable:

Review Comment:
   Same comment as 
[_SequenceRowIterable](https://github.com/apache/flink/pull/28797/changes#r3633002727)
 except this class also doesn't appear to be used at all



##########
flink-python/pyflink/dataframe/convert.py:
##########
@@ -0,0 +1,232 @@
+################################################################################
+#  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.
+################################################################################
+
+from typing import Any, Iterable, Iterator, List, Mapping, Optional, Sequence, 
Tuple, Union, cast
+
+from pyflink.dataframe.context import get_or_create_table_environment
+from pyflink.dataframe.dataframe import DataFrame
+from pyflink.util.api_stability_decorators import PublicEvolving
+
+__all__ = ["from_dict", "from_records"]
+
+
+class _SequenceRowIterable:
+    """Produce tuple rows from row-oriented sequence records."""
+
+    def __init__(self, records: Sequence[Sequence[Any]]):
+        self._records = records
+
+    def __iter__(self) -> Iterator[Sequence[Any]]:
+        return (tuple(record) for record in self._records)
+
+
+class _ColumnRowIterable:
+    """Produce tuple rows from column-oriented sequences."""
+
+    def __init__(
+        self, data: Mapping[str, Sequence[Any]], columns: Sequence[str]
+    ):
+        self._data = data
+        self._columns = columns
+
+    def __iter__(self) -> Iterator[Sequence[Any]]:
+        return zip(*(self._data[name] for name in self._columns))
+
+
+def _validate_schema(schema: List[str]) -> None:
+    if not isinstance(schema, list) or any(not isinstance(name, str) for name 
in schema):
+        raise TypeError("schema must be a list of strings")
+    if not schema:
+        raise ValueError("schema must not be empty")
+    if any(not name for name in schema):
+        raise ValueError("schema field names must not be empty")
+    if len(set(schema)) != len(schema):
+        raise ValueError("schema field names must be unique")
+
+
+def _is_named_tuple_record(record: Any) -> bool:
+    return isinstance(record, tuple) and isinstance(
+        getattr(record, "_fields", None), tuple
+    )
+
+
+def _to_field_mapping(
+    record: Any, named_tuple_records: bool

Review Comment:
   Is the purpose of the `named_tuple_records` boolean enforcement of 
standardization? If so I feel like a name change would help clarify things. 
Something like:
   
   ```python
   def _to_field_mapping(record: Any, enforce_named_tuple: bool):
       ...
   ```
   
   without seeing usage this kinda reads like redundant operations as is.



##########
flink-python/pyflink/dataframe/convert.py:
##########
@@ -0,0 +1,232 @@
+################################################################################
+#  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.
+################################################################################
+
+from typing import Any, Iterable, Iterator, List, Mapping, Optional, Sequence, 
Tuple, Union, cast
+
+from pyflink.dataframe.context import get_or_create_table_environment
+from pyflink.dataframe.dataframe import DataFrame
+from pyflink.util.api_stability_decorators import PublicEvolving
+
+__all__ = ["from_dict", "from_records"]
+
+
+class _SequenceRowIterable:

Review Comment:
   Why are we creating custom iterables? This class looks like it could be 
replaced with a simple:
   
   ```python
   rows = (tuple(record) for record in sequence_rows)
   ```
   
   on line 178



##########
flink-python/pyflink/dataframe/convert.py:
##########
@@ -0,0 +1,232 @@
+################################################################################
+#  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.
+################################################################################
+
+from typing import Any, Iterable, Iterator, List, Mapping, Optional, Sequence, 
Tuple, Union, cast
+
+from pyflink.dataframe.context import get_or_create_table_environment
+from pyflink.dataframe.dataframe import DataFrame
+from pyflink.util.api_stability_decorators import PublicEvolving
+
+__all__ = ["from_dict", "from_records"]
+
+
+class _SequenceRowIterable:
+    """Produce tuple rows from row-oriented sequence records."""
+
+    def __init__(self, records: Sequence[Sequence[Any]]):
+        self._records = records
+
+    def __iter__(self) -> Iterator[Sequence[Any]]:
+        return (tuple(record) for record in self._records)
+
+
+class _ColumnRowIterable:
+    """Produce tuple rows from column-oriented sequences."""
+
+    def __init__(
+        self, data: Mapping[str, Sequence[Any]], columns: Sequence[str]
+    ):
+        self._data = data
+        self._columns = columns
+
+    def __iter__(self) -> Iterator[Sequence[Any]]:
+        return zip(*(self._data[name] for name in self._columns))
+
+
+def _validate_schema(schema: List[str]) -> None:
+    if not isinstance(schema, list) or any(not isinstance(name, str) for name 
in schema):
+        raise TypeError("schema must be a list of strings")
+    if not schema:
+        raise ValueError("schema must not be empty")
+    if any(not name for name in schema):
+        raise ValueError("schema field names must not be empty")
+    if len(set(schema)) != len(schema):
+        raise ValueError("schema field names must be unique")
+
+
+def _is_named_tuple_record(record: Any) -> bool:
+    return isinstance(record, tuple) and isinstance(
+        getattr(record, "_fields", None), tuple
+    )
+
+
+def _to_field_mapping(
+    record: Any, named_tuple_records: bool
+) -> Mapping[str, Any]:
+    if named_tuple_records:
+        if not _is_named_tuple_record(record):
+            raise TypeError("records must all be named tuples")
+        fields = cast(Tuple[str, ...], getattr(record, "_fields"))
+        return dict(zip(fields, record))
+    if not isinstance(record, Mapping):
+        raise TypeError("records must all be mappings")
+    return cast(Mapping[str, Any], record)
+
+
+def _from_rows(rows: Iterable[Sequence[Any]], columns: Sequence[str]) -> 
DataFrame:
+    table = get_or_create_table_environment().from_elements(rows, 
list(columns))
+    return DataFrame(table)
+
+
+@PublicEvolving()
+def from_records(
+    data: Sequence[Union[Sequence[Any], Mapping[str, Any]]],
+    schema: Optional[List[str]] = None,
+) -> DataFrame:
+    """
+    Create a DataFrame from row-oriented records.
+
+    For mapping and named tuple records with an explicit ``schema``, every 
record must contain all
+    schema fields; other fields are ignored. When ``schema`` is omitted, the 
keys or fields from
+    the first record are used as the schema and every record must have exactly 
those fields.
+
+    For other sequence records, every record must have the same number of 
values. A ``schema`` is
+    required to provide the field names.
+
+    Field types are inferred from the record values.
+
+    :param data: Non-empty sequence of mapping or sequence records.
+    :param schema: Optional non-empty list of field names.
+    :return: A DataFrame containing the records.
+    :raises TypeError: If a record or schema has an invalid type.
+    :raises ValueError: If data or schema is empty, schema field names are 
invalid, a required
+        schema is omitted, a required field is absent, inferred record fields 
differ, or record
+        widths differ.
+
+    Example::
+
+        >>> import pyflink.dataframe as pf
+        >>> users = pf.from_records([
+        ...     {"id": 1, "name": "Alice"},
+        ...     {"id": 2, "name": "Bob"},
+        ... ])
+        >>> users = pf.from_records(
+        ...     [(1, "Alice"), (2, "Bob")], schema=["id", "name"]
+        ... )
+        >>> from typing import NamedTuple
+        >>> class User(NamedTuple):
+        ...     id: int
+        ...     name: str
+        >>> users = pf.from_records([User(1, "Alice"), User(2, "Bob")])
+        >>> selected_users = pf.from_records(
+        ...     [User(1, "Alice")], schema=["name", "id"]
+        ... )
+
+    .. versionadded:: 2.4.0
+    """
+    if not isinstance(data, Sequence):
+        raise TypeError("data must be a sequence of records")
+    if not data:
+        raise ValueError("data must not be empty")
+
+    first_record = data[0]
+    rows: Iterable[Sequence[Any]]
+    named_tuple_records = _is_named_tuple_record(first_record)
+    if isinstance(first_record, Mapping) or named_tuple_records:
+        first_fields = _to_field_mapping(first_record, named_tuple_records)
+        columns = list(first_fields.keys()) if schema is None else schema
+        _validate_schema(columns)
+        expected_fields = set(columns)
+        field_rows: List[Sequence[Any]] = []
+        for index, record in enumerate(data):
+            field_values = _to_field_mapping(record, named_tuple_records)
+            if schema is None and set(field_values.keys()) != expected_fields:

Review Comment:
   we have already validated that `columns` has unique fields. Do we expect the 
fields to have identical orderings? If so I think we can get rid of 
`expected_fields` and just do:
   
   ```python
   list(field_values.keys()) != columns
   ```



##########
flink-python/pyflink/dataframe/convert.py:
##########
@@ -0,0 +1,232 @@
+################################################################################
+#  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.
+################################################################################
+
+from typing import Any, Iterable, Iterator, List, Mapping, Optional, Sequence, 
Tuple, Union, cast
+
+from pyflink.dataframe.context import get_or_create_table_environment
+from pyflink.dataframe.dataframe import DataFrame
+from pyflink.util.api_stability_decorators import PublicEvolving
+
+__all__ = ["from_dict", "from_records"]
+
+
+class _SequenceRowIterable:
+    """Produce tuple rows from row-oriented sequence records."""
+
+    def __init__(self, records: Sequence[Sequence[Any]]):
+        self._records = records
+
+    def __iter__(self) -> Iterator[Sequence[Any]]:
+        return (tuple(record) for record in self._records)
+
+
+class _ColumnRowIterable:
+    """Produce tuple rows from column-oriented sequences."""
+
+    def __init__(
+        self, data: Mapping[str, Sequence[Any]], columns: Sequence[str]
+    ):
+        self._data = data
+        self._columns = columns
+
+    def __iter__(self) -> Iterator[Sequence[Any]]:
+        return zip(*(self._data[name] for name in self._columns))
+
+
+def _validate_schema(schema: List[str]) -> None:
+    if not isinstance(schema, list) or any(not isinstance(name, str) for name 
in schema):
+        raise TypeError("schema must be a list of strings")
+    if not schema:
+        raise ValueError("schema must not be empty")
+    if any(not name for name in schema):
+        raise ValueError("schema field names must not be empty")
+    if len(set(schema)) != len(schema):
+        raise ValueError("schema field names must be unique")
+
+
+def _is_named_tuple_record(record: Any) -> bool:
+    return isinstance(record, tuple) and isinstance(
+        getattr(record, "_fields", None), tuple
+    )
+
+
+def _to_field_mapping(
+    record: Any, named_tuple_records: bool
+) -> Mapping[str, Any]:
+    if named_tuple_records:
+        if not _is_named_tuple_record(record):
+            raise TypeError("records must all be named tuples")
+        fields = cast(Tuple[str, ...], getattr(record, "_fields"))
+        return dict(zip(fields, record))
+    if not isinstance(record, Mapping):
+        raise TypeError("records must all be mappings")
+    return cast(Mapping[str, Any], record)
+
+
+def _from_rows(rows: Iterable[Sequence[Any]], columns: Sequence[str]) -> 
DataFrame:
+    table = get_or_create_table_environment().from_elements(rows, 
list(columns))
+    return DataFrame(table)
+
+
+@PublicEvolving()
+def from_records(
+    data: Sequence[Union[Sequence[Any], Mapping[str, Any]]],
+    schema: Optional[List[str]] = None,
+) -> DataFrame:
+    """
+    Create a DataFrame from row-oriented records.
+
+    For mapping and named tuple records with an explicit ``schema``, every 
record must contain all
+    schema fields; other fields are ignored. When ``schema`` is omitted, the 
keys or fields from
+    the first record are used as the schema and every record must have exactly 
those fields.
+
+    For other sequence records, every record must have the same number of 
values. A ``schema`` is
+    required to provide the field names.
+
+    Field types are inferred from the record values.
+
+    :param data: Non-empty sequence of mapping or sequence records.
+    :param schema: Optional non-empty list of field names.
+    :return: A DataFrame containing the records.
+    :raises TypeError: If a record or schema has an invalid type.
+    :raises ValueError: If data or schema is empty, schema field names are 
invalid, a required
+        schema is omitted, a required field is absent, inferred record fields 
differ, or record
+        widths differ.
+
+    Example::
+
+        >>> import pyflink.dataframe as pf
+        >>> users = pf.from_records([
+        ...     {"id": 1, "name": "Alice"},
+        ...     {"id": 2, "name": "Bob"},
+        ... ])
+        >>> users = pf.from_records(
+        ...     [(1, "Alice"), (2, "Bob")], schema=["id", "name"]
+        ... )
+        >>> from typing import NamedTuple
+        >>> class User(NamedTuple):
+        ...     id: int
+        ...     name: str
+        >>> users = pf.from_records([User(1, "Alice"), User(2, "Bob")])
+        >>> selected_users = pf.from_records(
+        ...     [User(1, "Alice")], schema=["name", "id"]
+        ... )
+
+    .. versionadded:: 2.4.0
+    """
+    if not isinstance(data, Sequence):
+        raise TypeError("data must be a sequence of records")
+    if not data:
+        raise ValueError("data must not be empty")
+
+    first_record = data[0]
+    rows: Iterable[Sequence[Any]]
+    named_tuple_records = _is_named_tuple_record(first_record)

Review Comment:
   nit: can we change the name from `named_tuple_records` to 
`is_named_tuple_record`? The current naming implies a collection of tuple 
records rather than a boolean



##########
flink-python/pyflink/dataframe/context.py:
##########
@@ -0,0 +1,104 @@
+################################################################################
+#  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.
+################################################################################
+
+from typing import Optional, TYPE_CHECKING
+
+from pyflink.util.api_stability_decorators import PublicEvolving
+
+__all__ = [
+    "set_table_environment",
+    "get_table_environment",
+    "get_or_create_table_environment",
+]
+
+if TYPE_CHECKING:
+    from pyflink.table import StreamTableEnvironment
+
+
+_global_table_environment: Optional["StreamTableEnvironment"] = None

Review Comment:
   Might be slightly out of scope, but thoughts on making this a 
[`ContextVar`](https://docs.python.org/3/library/contextvars.html)? It wouldn't 
make much difference with the current planned API, but with the addition of a 
function like:
   
   ```python
   from contextlib import contextmanager
   
   @contextmanager
   def table_environment(t_env):
       global _global_table_environment
       with _global_table_environment.set(t_env):
           yield _global_table_environment.get()
   ```
   
   users would be able to selectively change their environments without 
worrying about mutating global state:
   
   ```python
   my_environment = ...
   set_table_environment(my_environment)
   assert get_table_environment() == my_environment
   
   # Start a temp environment
   temp_environment = ...
   with table_environment(temp_environment):
       assert get_table_environment() == temp_environment
   
   assert get_table_environment() == my_environment
   ```
   
   this would better support async or multithreaded environments



##########
flink-python/pyflink/dataframe/convert.py:
##########
@@ -0,0 +1,232 @@
+################################################################################
+#  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.
+################################################################################
+
+from typing import Any, Iterable, Iterator, List, Mapping, Optional, Sequence, 
Tuple, Union, cast
+
+from pyflink.dataframe.context import get_or_create_table_environment
+from pyflink.dataframe.dataframe import DataFrame
+from pyflink.util.api_stability_decorators import PublicEvolving
+
+__all__ = ["from_dict", "from_records"]
+
+
+class _SequenceRowIterable:
+    """Produce tuple rows from row-oriented sequence records."""
+
+    def __init__(self, records: Sequence[Sequence[Any]]):
+        self._records = records
+
+    def __iter__(self) -> Iterator[Sequence[Any]]:
+        return (tuple(record) for record in self._records)
+
+
+class _ColumnRowIterable:
+    """Produce tuple rows from column-oriented sequences."""
+
+    def __init__(
+        self, data: Mapping[str, Sequence[Any]], columns: Sequence[str]
+    ):
+        self._data = data
+        self._columns = columns
+
+    def __iter__(self) -> Iterator[Sequence[Any]]:
+        return zip(*(self._data[name] for name in self._columns))
+
+
+def _validate_schema(schema: List[str]) -> None:
+    if not isinstance(schema, list) or any(not isinstance(name, str) for name 
in schema):
+        raise TypeError("schema must be a list of strings")
+    if not schema:
+        raise ValueError("schema must not be empty")
+    if any(not name for name in schema):
+        raise ValueError("schema field names must not be empty")
+    if len(set(schema)) != len(schema):
+        raise ValueError("schema field names must be unique")
+
+
+def _is_named_tuple_record(record: Any) -> bool:
+    return isinstance(record, tuple) and isinstance(
+        getattr(record, "_fields", None), tuple
+    )
+
+
+def _to_field_mapping(
+    record: Any, named_tuple_records: bool

Review Comment:
   I feel like the type of record could be described in a more readable fashion 
with an enum:
   
   ```python
   from enum import Enum
   
   class RecordType(Enum)
       NamedTuple = 'named_tuple'
       Mapping = 'mapping'
       Sequence = 'sequence'
   
       @classmethod
       def from_record(cls, record: Any) -> 'RecordType':
           if isinstance(record, tuple) and isinstance(getattr(record, 
"_fields", None), tuple):
               return cls.NamedTuple
           elif isinstance(record, Mapping):
               return cls.Mapping
           elif isinstance(record, Sequence):
               return cls.Sequence
           else:
               raise TypeError("records must be mappings or sequences")
   ```
   
   this makes the supported types more obvious and encourages variable names 
that I think would better encompass intent.
   
   Current usage could look like:
   
   ```python
   def from_records(
       data: Sequence[Union[Sequence[Any], Mapping[str, Any]]],
       schema: Optional[List[str]] = None,
   ) -> DataFrame:
       
       ...
       
       expected_record_type = RecordType.from_record(first_record)
   
       ...
   
       if expected_record_type in [RecordType.Mapping, RecordType.NamedTuple]:
           first_fields = _to_field_mapping(first_record, expected_record_type)
           ...
       else: # We know it must be a RecordType.Sequence since we are enforcing 
this on creation
           ...
   ```
   
   ```python
   def _to_field_mapping(record: Any, expected_type: RecordType) -> 
Mapping[str, Any]:
       record_type = RecordType.from_record(record)
       if expected_type == RecordType.NamedTuple:
           if record_type != expected_type:
               raise TypeError("records must all be named tuples")
           ...
       elif record_type != RecordType.Mapping:
           raise TypeError("records must all be mappings")
       # We should probably handle a sequence type here for completeness
       return cast(Mapping[str, Any], record)
   ```
   
   looking forward this would also work really well with structural pattern 
matching once the lib drops support for python 3.9:
   
   ```python
   match RecordType.from_record(record):
       case RecordType.NamedTuple | RecordType.Mapping:
          # Do something
       case RecordType.Sequence:
          # Do something else
       case _:
           raise TypeError('unsupported record type')
   ```



-- 
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]

Reply via email to