mirror of
https://github.com/datahub-project/datahub.git
synced 2025-09-01 05:13:15 +00:00
fix(ingestion/airflow-plugin): warning log for non-materialized iolets (#10421)
This commit is contained in:
parent
de5b503380
commit
4d20c22a49
@ -428,8 +428,8 @@ class AirflowGenerator:
|
|||||||
dpi = DataProcessInstance.from_datajob(
|
dpi = DataProcessInstance.from_datajob(
|
||||||
datajob=datajob,
|
datajob=datajob,
|
||||||
id=f"{dag.dag_id}_{ti.task_id}_{dag_run.run_id}",
|
id=f"{dag.dag_id}_{ti.task_id}_{dag_run.run_id}",
|
||||||
clone_inlets=True,
|
clone_inlets=config is None or config.materialize_iolets,
|
||||||
clone_outlets=True,
|
clone_outlets=config is None or config.materialize_iolets,
|
||||||
)
|
)
|
||||||
job_property_bag: Dict[str, str] = {}
|
job_property_bag: Dict[str, str] = {}
|
||||||
job_property_bag["run_id"] = str(dag_run.run_id)
|
job_property_bag["run_id"] = str(dag_run.run_id)
|
||||||
|
@ -433,6 +433,14 @@ class DataHubListener:
|
|||||||
|
|
||||||
self.emitter.emit(operation_mcp)
|
self.emitter.emit(operation_mcp)
|
||||||
logger.debug(f"Emitted Dataset Operation: {outlet}")
|
logger.debug(f"Emitted Dataset Operation: {outlet}")
|
||||||
|
else:
|
||||||
|
if self.graph:
|
||||||
|
for outlet in datajob.outlets:
|
||||||
|
if not self.graph.exists(str(outlet)):
|
||||||
|
logger.warning(f"Dataset {str(outlet)} not materialized")
|
||||||
|
for inlet in datajob.inlets:
|
||||||
|
if not self.graph.exists(str(inlet)):
|
||||||
|
logger.warning(f"Dataset {str(inlet)} not materialized")
|
||||||
|
|
||||||
def on_task_instance_finish(
|
def on_task_instance_finish(
|
||||||
self, task_instance: "TaskInstance", status: InstanceRunResult
|
self, task_instance: "TaskInstance", status: InstanceRunResult
|
||||||
|
Loading…
x
Reference in New Issue
Block a user