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]