justinpakzad commented on code in PR #70446:
URL: https://github.com/apache/airflow/pull/70446#discussion_r4112057546
##########
airflow-core/src/airflow/cli/commands/dag_command.py:
##########
@@ -80,6 +82,95 @@
_RUN_CHUNK_SIZE = 500
+def _normalize_serialized_dag_for_stability_check(serialized_dag: dict[str,
Any]) -> dict[str, Any]:
+ normalized = SerializedDagModel._sort_serialized_dag_dict(serialized_dag)
+ normalized["dag"].pop("fileloc", None)
+ normalized["dag"].pop("bundle_name", None)
+ return normalized
+
+
+def _serialize_dag_for_stability_check(dag: DAG) -> tuple[str, dict[str, Any]]:
+ serialized_dag = DagSerialization.to_dict(dag)
+ return SerializedDagModel.hash(serialized_dag),
_normalize_serialized_dag_for_stability_check(
+ serialized_dag
+ )
+
+
+def _format_stability_diff(
+ dag_id: str,
+ first_serialized_dag: dict[str, Any],
+ second_serialized_dag: dict[str, Any],
+) -> str:
+ before = json.dumps(first_serialized_dag, indent=2,
sort_keys=True).splitlines()
+ after = json.dumps(second_serialized_dag, indent=2,
sort_keys=True).splitlines()
+ diff = difflib.unified_diff(
+ before,
+ after,
+ fromfile=f"parse 1: {dag_id} ",
+ tofile=f"parse 2: {dag_id} ",
+ lineterm="",
+ )
+ return "\n".join(diff)
+
+
+def _parse_dags_for_stability_check(dag_folder: str | None) -> DagBag:
+ return DagBag(dag_folder=dag_folder, load_op_links=False)
+
+
+@cli_utils.action_cli
+@providers_configuration_loaded
+def dag_stability_check(args) -> None:
+ dag_hashes_by_id: dict[str, list[str]] = {}
+ serialized_dags_by_id: dict[str, list[dict[str, Any]]] = {}
+ seen_dag_ids: set[str] = set()
+ nParse = (
Review Comment:
Nit: I think co-pilot noted this already but we should be using snake case.
```suggestion
n_parse = (
```
##########
airflow-core/src/airflow/cli/commands/dag_command.py:
##########
@@ -80,6 +82,95 @@
_RUN_CHUNK_SIZE = 500
+def _normalize_serialized_dag_for_stability_check(serialized_dag: dict[str,
Any]) -> dict[str, Any]:
+ normalized = SerializedDagModel._sort_serialized_dag_dict(serialized_dag)
+ normalized["dag"].pop("fileloc", None)
+ normalized["dag"].pop("bundle_name", None)
+ return normalized
+
+
+def _serialize_dag_for_stability_check(dag: DAG) -> tuple[str, dict[str, Any]]:
+ serialized_dag = DagSerialization.to_dict(dag)
+ return SerializedDagModel.hash(serialized_dag),
_normalize_serialized_dag_for_stability_check(
+ serialized_dag
+ )
+
+
+def _format_stability_diff(
+ dag_id: str,
+ first_serialized_dag: dict[str, Any],
+ second_serialized_dag: dict[str, Any],
+) -> str:
+ before = json.dumps(first_serialized_dag, indent=2,
sort_keys=True).splitlines()
+ after = json.dumps(second_serialized_dag, indent=2,
sort_keys=True).splitlines()
+ diff = difflib.unified_diff(
+ before,
+ after,
+ fromfile=f"parse 1: {dag_id} ",
+ tofile=f"parse 2: {dag_id} ",
+ lineterm="",
+ )
+ return "\n".join(diff)
+
+
+def _parse_dags_for_stability_check(dag_folder: str | None) -> DagBag:
+ return DagBag(dag_folder=dag_folder, load_op_links=False)
+
+
+@cli_utils.action_cli
+@providers_configuration_loaded
+def dag_stability_check(args) -> None:
+ dag_hashes_by_id: dict[str, list[str]] = {}
+ serialized_dags_by_id: dict[str, list[dict[str, Any]]] = {}
+ seen_dag_ids: set[str] = set()
+ nParse = (
+ 2 # NOTE: Set parsing number 2, in most of the case twice parse
should catch the stability issues.
+ )
+
+ for _iteration in range(1, nParse + 1):
Review Comment:
Nit: we should use `_` for placeholder variables that aren't used inside the
loop at all.
```suggestion
for _ in range(1, nParse + 1):
```
##########
airflow-core/src/airflow/cli/commands/dag_command.py:
##########
@@ -80,6 +82,95 @@
_RUN_CHUNK_SIZE = 500
+def _normalize_serialized_dag_for_stability_check(serialized_dag: dict[str,
Any]) -> dict[str, Any]:
+ normalized = SerializedDagModel._sort_serialized_dag_dict(serialized_dag)
+ normalized["dag"].pop("fileloc", None)
+ normalized["dag"].pop("bundle_name", None)
+ return normalized
+
+
+def _serialize_dag_for_stability_check(dag: DAG) -> tuple[str, dict[str, Any]]:
+ serialized_dag = DagSerialization.to_dict(dag)
+ return SerializedDagModel.hash(serialized_dag),
_normalize_serialized_dag_for_stability_check(
+ serialized_dag
+ )
+
+
+def _format_stability_diff(
+ dag_id: str,
+ first_serialized_dag: dict[str, Any],
+ second_serialized_dag: dict[str, Any],
+) -> str:
+ before = json.dumps(first_serialized_dag, indent=2,
sort_keys=True).splitlines()
+ after = json.dumps(second_serialized_dag, indent=2,
sort_keys=True).splitlines()
+ diff = difflib.unified_diff(
+ before,
+ after,
+ fromfile=f"parse 1: {dag_id} ",
+ tofile=f"parse 2: {dag_id} ",
+ lineterm="",
+ )
+ return "\n".join(diff)
+
+
+def _parse_dags_for_stability_check(dag_folder: str | None) -> DagBag:
+ return DagBag(dag_folder=dag_folder, load_op_links=False)
+
+
+@cli_utils.action_cli
+@providers_configuration_loaded
+def dag_stability_check(args) -> None:
+ dag_hashes_by_id: dict[str, list[str]] = {}
+ serialized_dags_by_id: dict[str, list[dict[str, Any]]] = {}
+ seen_dag_ids: set[str] = set()
+ nParse = (
+ 2 # NOTE: Set parsing number 2, in most of the case twice parse
should catch the stability issues.
+ )
+
+ for _iteration in range(1, nParse + 1):
+ dagbag = _parse_dags_for_stability_check(args.dag_folder)
+
+ dags = dagbag.dags
+ if args.dag_id is not None:
+ dags = {args.dag_id: dagbag.dags[args.dag_id]} if args.dag_id in
dagbag.dags else {}
+
+ for _dag_id, dag in sorted(dags.items()):
+ dag_hash, serialized_dag = _serialize_dag_for_stability_check(dag)
+ dag_hashes_by_id.setdefault(_dag_id, []).append(dag_hash)
+ serialized_dags_by_id.setdefault(_dag_id,
[]).append(serialized_dag)
+ seen_dag_ids.add(_dag_id)
+
+ if args.fail_fast:
+ for _dag_id, dag_hashes in dag_hashes_by_id.items():
+ if len(set(dag_hashes)) > 1:
+ break
+ else:
+ continue
+ break
Review Comment:
I think this is a no-op if the number of parses is fixed at 2. On the first
iteration, each Dag will only have 1 hash, `if len(set(dag_hashes)) > 1:` can
only be True on the second iteration, which is already the last iteration of
`range(1, nParse + 1)` so the loop ends regardless.
Related to my other comment, if we made the number of parses something a
user can pass in (e.g. 5), then I think this would be useful.
##########
airflow-core/src/airflow/cli/commands/dag_command.py:
##########
@@ -80,6 +82,95 @@
_RUN_CHUNK_SIZE = 500
+def _normalize_serialized_dag_for_stability_check(serialized_dag: dict[str,
Any]) -> dict[str, Any]:
+ normalized = SerializedDagModel._sort_serialized_dag_dict(serialized_dag)
+ normalized["dag"].pop("fileloc", None)
+ normalized["dag"].pop("bundle_name", None)
+ return normalized
+
+
+def _serialize_dag_for_stability_check(dag: DAG) -> tuple[str, dict[str, Any]]:
+ serialized_dag = DagSerialization.to_dict(dag)
+ return SerializedDagModel.hash(serialized_dag),
_normalize_serialized_dag_for_stability_check(
+ serialized_dag
+ )
+
+
+def _format_stability_diff(
+ dag_id: str,
+ first_serialized_dag: dict[str, Any],
+ second_serialized_dag: dict[str, Any],
+) -> str:
+ before = json.dumps(first_serialized_dag, indent=2,
sort_keys=True).splitlines()
+ after = json.dumps(second_serialized_dag, indent=2,
sort_keys=True).splitlines()
+ diff = difflib.unified_diff(
+ before,
+ after,
+ fromfile=f"parse 1: {dag_id} ",
+ tofile=f"parse 2: {dag_id} ",
+ lineterm="",
+ )
+ return "\n".join(diff)
+
+
+def _parse_dags_for_stability_check(dag_folder: str | None) -> DagBag:
+ return DagBag(dag_folder=dag_folder, load_op_links=False)
+
+
+@cli_utils.action_cli
+@providers_configuration_loaded
+def dag_stability_check(args) -> None:
+ dag_hashes_by_id: dict[str, list[str]] = {}
+ serialized_dags_by_id: dict[str, list[dict[str, Any]]] = {}
+ seen_dag_ids: set[str] = set()
+ nParse = (
+ 2 # NOTE: Set parsing number 2, in most of the case twice parse
should catch the stability issues.
+ )
+
+ for _iteration in range(1, nParse + 1):
+ dagbag = _parse_dags_for_stability_check(args.dag_folder)
Review Comment:
If the folder provided by the args doesn't exist the command would succeed
and return "Dag stability check passed for 0 Dags(s)". Should a non-existent
path fail here, given a missing --dag-id does?
##########
airflow-core/src/airflow/cli/commands/dag_command.py:
##########
@@ -80,6 +82,95 @@
_RUN_CHUNK_SIZE = 500
+def _normalize_serialized_dag_for_stability_check(serialized_dag: dict[str,
Any]) -> dict[str, Any]:
+ normalized = SerializedDagModel._sort_serialized_dag_dict(serialized_dag)
+ normalized["dag"].pop("fileloc", None)
+ normalized["dag"].pop("bundle_name", None)
+ return normalized
+
+
+def _serialize_dag_for_stability_check(dag: DAG) -> tuple[str, dict[str, Any]]:
+ serialized_dag = DagSerialization.to_dict(dag)
+ return SerializedDagModel.hash(serialized_dag),
_normalize_serialized_dag_for_stability_check(
+ serialized_dag
+ )
+
+
+def _format_stability_diff(
+ dag_id: str,
+ first_serialized_dag: dict[str, Any],
+ second_serialized_dag: dict[str, Any],
+) -> str:
+ before = json.dumps(first_serialized_dag, indent=2,
sort_keys=True).splitlines()
+ after = json.dumps(second_serialized_dag, indent=2,
sort_keys=True).splitlines()
+ diff = difflib.unified_diff(
+ before,
+ after,
+ fromfile=f"parse 1: {dag_id} ",
+ tofile=f"parse 2: {dag_id} ",
+ lineterm="",
+ )
+ return "\n".join(diff)
+
+
+def _parse_dags_for_stability_check(dag_folder: str | None) -> DagBag:
+ return DagBag(dag_folder=dag_folder, load_op_links=False)
+
+
+@cli_utils.action_cli
+@providers_configuration_loaded
+def dag_stability_check(args) -> None:
+ dag_hashes_by_id: dict[str, list[str]] = {}
+ serialized_dags_by_id: dict[str, list[dict[str, Any]]] = {}
+ seen_dag_ids: set[str] = set()
+ nParse = (
+ 2 # NOTE: Set parsing number 2, in most of the case twice parse
should catch the stability issues.
Review Comment:
Do you think we should make this a configurable value with a default of two
(e.g., an optional flag --n_parses)?
##########
airflow-core/src/airflow/cli/commands/dag_command.py:
##########
@@ -80,6 +82,95 @@
_RUN_CHUNK_SIZE = 500
+def _normalize_serialized_dag_for_stability_check(serialized_dag: dict[str,
Any]) -> dict[str, Any]:
+ normalized = SerializedDagModel._sort_serialized_dag_dict(serialized_dag)
+ normalized["dag"].pop("fileloc", None)
+ normalized["dag"].pop("bundle_name", None)
Review Comment:
Any reason these are popped here but still included in the `hash()` used for
comparison?
--
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]