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


##########
flink-python/pyflink/dataframe/convert.py:
##########
@@ -120,10 +148,239 @@ def _validate_schema(schema: List[str]) -> None:
         raise ValueError("schema field names must be unique")
 
 
+def _resolve_column_names(
+    input_names: Sequence[str], schema: Optional[List[str]]
+) -> List[str]:
+    column_names = list(input_names) if schema is None else schema
+    if (
+        schema is not None
+        and isinstance(schema, list)
+        and len(schema) != len(input_names)
+    ):
+        raise ValueError(
+            f"schema has {len(schema)} fields but data has "
+            f"{len(input_names)} columns"
+        )
+    _validate_schema(column_names)
+    return column_names
+
+
+def _parse_watermark(

Review Comment:
   Can we move this under the `_WaterMarkSpec` class? Like:
   
   ```python
   class Watermark:
   
   
        @classmethod
        def parse(cls, ....):
            ....
            return cls(....)
   ```



##########
flink-python/pyflink/dataframe/convert.py:
##########
@@ -120,10 +148,239 @@ def _validate_schema(schema: List[str]) -> None:
         raise ValueError("schema field names must be unique")
 
 
+def _resolve_column_names(
+    input_names: Sequence[str], schema: Optional[List[str]]
+) -> List[str]:
+    column_names = list(input_names) if schema is None else schema
+    if (
+        schema is not None
+        and isinstance(schema, list)
+        and len(schema) != len(input_names)
+    ):
+        raise ValueError(
+            f"schema has {len(schema)} fields but data has "
+            f"{len(input_names)} columns"
+        )
+    _validate_schema(column_names)
+    return column_names
+
+
+def _parse_watermark(
+    watermark: Optional[Tuple[str, str]],
+) -> Optional[_WatermarkSpec]:
+    if watermark is None:
+        return None
+    if not isinstance(watermark, tuple) or len(watermark) != 2:
+        raise TypeError("watermark must be a tuple of (column, expression)")
+    if any(not isinstance(value, str) or not value.strip() for value in 
watermark):
+        raise TypeError("watermark column and expression must be non-empty 
strings")
+    return _WatermarkSpec(*watermark)
+
+
+def _normalize_watermark_row_type(row_type: RowType, watermark: 
_WatermarkSpec) -> RowType:

Review Comment:
   Can we move this under the `_WaterMarkSpec` class? Like:
   
   ```python
   class Watermark:
   
       def normalize_row_type(self, row_type: RowType) -> RowType:
           ...
   ```



##########
flink-python/pyflink/dataframe/convert.py:
##########
@@ -120,10 +148,239 @@ def _validate_schema(schema: List[str]) -> None:
         raise ValueError("schema field names must be unique")
 
 
+def _resolve_column_names(
+    input_names: Sequence[str], schema: Optional[List[str]]
+) -> List[str]:
+    column_names = list(input_names) if schema is None else schema
+    if (
+        schema is not None
+        and isinstance(schema, list)
+        and len(schema) != len(input_names)
+    ):
+        raise ValueError(
+            f"schema has {len(schema)} fields but data has "
+            f"{len(input_names)} columns"
+        )
+    _validate_schema(column_names)
+    return column_names
+
+
+def _parse_watermark(
+    watermark: Optional[Tuple[str, str]],
+) -> Optional[_WatermarkSpec]:
+    if watermark is None:
+        return None
+    if not isinstance(watermark, tuple) or len(watermark) != 2:
+        raise TypeError("watermark must be a tuple of (column, expression)")
+    if any(not isinstance(value, str) or not value.strip() for value in 
watermark):
+        raise TypeError("watermark column and expression must be non-empty 
strings")
+    return _WatermarkSpec(*watermark)
+
+
+def _normalize_watermark_row_type(row_type: RowType, watermark: 
_WatermarkSpec) -> RowType:
+    column_name = watermark.column
+    matching_fields = [field for field in row_type.fields if field.name == 
column_name]
+    if not matching_fields:
+        raise ValueError(f"watermark column {column_name!r} is not present in 
data")
+
+    watermark_type = matching_fields[0].data_type
+    if not isinstance(watermark_type, (TimestampType, 
LocalZonedTimestampType)):
+        raise ValueError(
+            f"watermark column {column_name!r} must have a timestamp type"
+        )
+
+    fields = []
+    for field in row_type.fields:
+        data_type = field.data_type
+        if field.name == column_name and data_type.precision != 3:
+            data_type = type(data_type)(3, data_type._nullable)
+        fields.append(RowField(field.name, data_type, field.description))
+    return RowType(fields, row_type._nullable)
+
+
+def _resolve_watermark_schema(
+    row_type: RowType, watermark: Optional[_WatermarkSpec]
+) -> Tuple[RowType, Optional[Schema]]:
+    if watermark is None:
+        return row_type, None
+
+    row_type = _normalize_watermark_row_type(row_type, watermark)
+    table_schema = (
+        Schema.new_builder()
+        .from_row_data_type(row_type)
+        .watermark(watermark.column, watermark.expression)

Review Comment:
   since the watermark is a `NamedTuple`, you can simplify the call by 
unpacking it like:
   
   ```python
   schema.watermark(*watermark)
   ```
   
   I feel like this both reads better and is more concise.



##########
flink-python/pyflink/dataframe/tests/test_convert.py:
##########


Review Comment:
   since `_WatermarkSpec` is a `NamedTuple` you should add a quick test to 
verify unpacking order. This would help guard against errors caused by someone 
changing the internal field order or values. In particular, this would enable 
us to be confident that:
   
   ```python
   schema.watermark(*watermark)
   ```
   
   won't quietly break in the future.



##########
flink-python/pyflink/dataframe/convert.py:
##########
@@ -120,10 +148,239 @@ def _validate_schema(schema: List[str]) -> None:
         raise ValueError("schema field names must be unique")
 
 
+def _resolve_column_names(
+    input_names: Sequence[str], schema: Optional[List[str]]
+) -> List[str]:
+    column_names = list(input_names) if schema is None else schema
+    if (
+        schema is not None
+        and isinstance(schema, list)
+        and len(schema) != len(input_names)
+    ):
+        raise ValueError(
+            f"schema has {len(schema)} fields but data has "
+            f"{len(input_names)} columns"
+        )
+    _validate_schema(column_names)
+    return column_names
+
+
+def _parse_watermark(
+    watermark: Optional[Tuple[str, str]],
+) -> Optional[_WatermarkSpec]:
+    if watermark is None:
+        return None
+    if not isinstance(watermark, tuple) or len(watermark) != 2:
+        raise TypeError("watermark must be a tuple of (column, expression)")
+    if any(not isinstance(value, str) or not value.strip() for value in 
watermark):
+        raise TypeError("watermark column and expression must be non-empty 
strings")
+    return _WatermarkSpec(*watermark)
+
+
+def _normalize_watermark_row_type(row_type: RowType, watermark: 
_WatermarkSpec) -> RowType:
+    column_name = watermark.column
+    matching_fields = [field for field in row_type.fields if field.name == 
column_name]
+    if not matching_fields:
+        raise ValueError(f"watermark column {column_name!r} is not present in 
data")
+
+    watermark_type = matching_fields[0].data_type
+    if not isinstance(watermark_type, (TimestampType, 
LocalZonedTimestampType)):
+        raise ValueError(
+            f"watermark column {column_name!r} must have a timestamp type"
+        )
+
+    fields = []
+    for field in row_type.fields:
+        data_type = field.data_type
+        if field.name == column_name and data_type.precision != 3:
+            data_type = type(data_type)(3, data_type._nullable)
+        fields.append(RowField(field.name, data_type, field.description))
+    return RowType(fields, row_type._nullable)
+
+
+def _resolve_watermark_schema(

Review Comment:
   This helper function feels shallow and with only 2 uses a little unecessary. 
Thoughts on inlining the logic at the call sites like:
   
   ```python
   if watermark:
       row_type = watermark.normalize_row_type(row_type)
       table_schema = ...
   else:
       table_schema = None
   ```
   
   this feels a easier to follow to me



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