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'}

More resources

Learn more