Python Microsoft Fabric - Pipelines
Data Factory pipelines started on demand with parameters, monitored and cancelled from Python services.
Data pipelines are Fabric's ETL and ELT engine - Data Factory workflows that copy data, run transformations and orchestrate other items. Triggering them from your services connects your integrations to your data platform - a file arrives over SFTP, a pipeline loads it into the lakehouse. You create a Fabric connection in the Dashboard and every pipeline is available to your services.
Finding pipelines
Pipelines are workspace items of the DataPipeline type.
# -*- coding: utf-8 -*-
# Zato
from zato.server.service import Service
class ListPipelines(Service):
input = 'workspace_id'
def handle(self):
# Get the connection by its Dashboard name
conn = self.microsoft.fabric['My Fabric']
# List only the pipeline items
response = conn.list_items(self.request.input.workspace_id, 'DataPipeline')
pipelines = []
for item in response['value']:
pipelines.append({
'id': item['id'],
'name': item['displayName'],
})
self.response.payload = {'pipelines': pipelines}
Running a pipeline on demand
conn.run_job with the Pipeline job type starts a pipeline run.
# -*- coding: utf-8 -*-
# Zato
from zato.server.service import Service
class LoadIncomingOrders(Service):
input = 'workspace_id', 'pipeline_id'
def handle(self):
conn = self.microsoft.fabric['My Fabric']
# Start the pipeline
conn.run_job(self.request.input.workspace_id, self.request.input.pipeline_id, 'Pipeline')
self.response.payload = {'status': 'started'}
Passing parameters
Pipeline parameters travel in the job's execution payload - the pipeline reads them like parameters from any other trigger, so one pipeline can serve many sources.
# -*- coding: utf-8 -*-
# Zato
from zato.server.service import Service
class LoadSourceFile(Service):
input = 'workspace_id', 'pipeline_id', 'file_name'
def handle(self):
conn = self.microsoft.fabric['My Fabric']
# Parameters the pipeline will receive
payload = {
'executionData': {
'parameters': {
'source_file': self.request.input.file_name,
'target_table': 'orders',
}
}
}
conn.run_job(self.request.input.workspace_id, self.request.input.pipeline_id, 'Pipeline', payload)
self.response.payload = {'status': 'started', 'file': self.request.input.file_name}
Monitoring a pipeline run
conn.get_job returns the run's status - NotStarted, InProgress, Completed, Failed or Cancelled - which is what failure reports and completion checks build on.
# -*- coding: utf-8 -*-
# Zato
from zato.server.service import Service
class CheckPipelineRun(Service):
input = 'workspace_id', 'pipeline_id', 'job_id'
def handle(self):
conn = self.microsoft.fabric['My Fabric']
job = conn.get_job(
self.request.input.workspace_id,
self.request.input.pipeline_id,
self.request.input.job_id,
)
# Let the on-call team know if the load failed
if job['status'] == 'Failed':
self.logger.warning('Pipeline run failed: %s', job['id'])
self.response.payload = {'status': job['status']}
Cancelling a run
conn.cancel_job stops a run that should no longer continue - for instance, when the file it was loading turns out to be corrupted.
# -*- coding: utf-8 -*-
# Zato
from zato.server.service import Service
class CancelPipelineRun(Service):
input = 'workspace_id', 'pipeline_id', 'job_id'
def handle(self):
conn = self.microsoft.fabric['My Fabric']
conn.cancel_job(
self.request.input.workspace_id,
self.request.input.pipeline_id,
self.request.input.job_id,
)
self.response.payload = {'status': 'cancelled'}