rich-sandwich-70692
07/19/2023, 10:56 AMcurved-planet-99787
07/19/2023, 11:19 AMrich-sandwich-70692
07/19/2023, 11:20 AMrich-sandwich-70692
07/19/2023, 11:36 AMcurved-planet-99787
07/19/2023, 11:44 AMrich-sandwich-70692
07/19/2023, 1:27 PMrich-sandwich-70692
07/19/2023, 1:53 PMcurved-planet-99787
07/19/2023, 2:03 PMDataProcessInstance with a link to the corresponding job?rich-sandwich-70692
07/19/2023, 2:05 PMimport time
import uuid
from datahub.api.entities.corpgroup.corpgroup import CorpGroup
from datahub.api.entities.corpuser.corpuser import CorpUser
from datahub.api.entities.datajob.dataflow import DataFlow
from datahub.api.entities.datajob.datajob import DataJob
from datahub.api.entities.dataprocess.dataprocess_instance import (
DataProcessInstance,
InstanceRunResult,
)
from datahub.emitter.rest_emitter import DatahubRestEmitter
emitter = DatahubRestEmitter("<http://localhost:8080>")
jobFlow = DataFlow(env="prod", orchestrator="airflow", id="flow2")
jobFlow.emit(emitter)
# Flowurn as constructor
dataJob = DataJob(flow_urn=jobFlow.urn, id="job1", name="My Job 1")
dataJob.properties["custom_properties"] = "test"
dataJob.emit(emitter)
jobFlowRun: DataProcessInstance = DataProcessInstance(
orchestrator="airflow", cluster="prod", id=f"{jobFlow.id}-{uuid.uuid4()}"
)
jobRun1: DataProcessInstance = DataProcessInstance(
orchestrator="airflow",
cluster="prod",
id=f"{jobFlow.id}-{dataJob.id}-{uuid.uuid4()}",
)
jobRun1.parent_instance = jobFlowRun.urn
jobRun1.template_urn = dataJob.urn
jobRun1.emit_process_start(
emitter=emitter, start_timestamp_millis=int(time.time() * 1000), emit_template=False
)
jobRun1.emit_process_end(
emitter=emitter,
end_timestamp_millis=int(time.time() * 1000),
result=InstanceRunResult.FAILURE,
)
Also attaching the job runs in the UIcurved-planet-99787
07/19/2023, 2:07 PMinlets and outlets, right?rich-sandwich-70692
07/19/2023, 2:09 PMDataProcessInstance
? I didn't find any reference to them here: https://datahubproject.io/docs/generated/metamodel/entities/dataprocessinstance/#dataprocessinstancerunevent-timeseriescurved-planet-99787
07/19/2023, 2:12 PMrich-sandwich-70692
07/19/2023, 2:16 PMoutlets=["urn:li:dataset:(urn:li:dataPlatform:clickhouse,GlobalLicenseServer.GlobalLicenseServer.license_usages,PROD)"]
?curved-planet-99787
07/19/2023, 2:21 PMDatasetUrn instances. So you could do the following:
outlets=[DatasetUrn.create_from_string("urn:li:dataset:(urn:li:dataPlatform:clickhouse,GlobalLicenseServer.GlobalLicenseServer.license_usages,PROD)")]rich-sandwich-70692
07/19/2023, 2:54 PMrich-sandwich-70692
07/19/2023, 2:56 PMjobFlowRun: DataProcessInstance = DataProcessInstance(
orchestrator="airflow", cluster="prod", id=f"{jobFlow.id}-{uuid.uuid4()}",
outlets=[DatasetUrn.create_from_string(
"urn:li:dataset:(urn:li:dataPlatform:clickhouse,GlobalLicenseServer.GlobalLicenseServer.license_usages,PROD)")]
)
And the outlets still don't show up in the job run and "Operations" does not appear in the UI of the table
I took the urn from the url of the table, perhaps that is not the way?rich-sandwich-70692
07/19/2023, 2:59 PMinlets (List[str]): List of entities the DataProcessInstance consumes
outlets (List[str]): List of entities the DataProcessInstance produces
Sorry for getting on your case so much š
rich-sandwich-70692
07/19/2023, 3:08 PMcurved-planet-99787
07/20/2023, 5:35 AMrich-sandwich-70692
07/20/2023, 9:14 AMcurved-planet-99787
07/20/2023, 10:44 AM