Skip to main content
Unlisted page
This page is unlisted. Search engines will not index it, and only users having a direct link can access it.

Lakeflow Integrations

Beta

This feature is in Beta.

Integrations are custom tasks that you can add to Lakeflow Jobs. An author writes an integration in Python and registers it in the workspace, making it available from the Add task dialog. Any other user can then use it in a job without writing code.

An integration is either a function or a sensor:

  • Functions perform a one-time operation, such as sending a notification.
  • Sensors wait on a condition by checking it in a loop. Between checks, a sensor releases its compute instead of holding it idle, so it can wait efficiently.

This page covers both adding an integration and using a registered integration.

note

To provide feedback or ask questions during the preview, email lakeflow-integrations-private-preview@databricks.com.

Add an integration​

You author integrations using Declarative Automation Bundles. Create a bundle from the Lakeflow Integrations template, which scaffolds two sample integrations that you can modify to build your own. You can create the bundle from the workspace or from the Databricks CLI.

To create and register an integration from the workspace UI:

  1. Create a bundle from the Lakeflow Integrations template, following Tutorial: Create and deploy a bundle in the workspace.

  2. When the bundle is created, click the bundle (rocket) icon in the left sidebar, then click Deploy. Deployment uploads the integration YAML and wheel files, and creates example jobs.

  3. Find the uploaded wheels and YAML files at /Workspace/Users/<user>/.bundle/<bundle>/dev/artifacts/.internal.

  4. Register the uploaded integrations so they appear in the UI, using a .lakeflow_integrations.yml file:

    • To make an integration available to any user, add it to /Workspace/.lakeflow_integrations.yml and grant workspace users permission to read the file.
    • To make an integration visible only to yourself, add it to /Workspace/Users/<user>/.lakeflow_integrations.yml.

    For example:

    YAML
    integrations:
    - '/Workspace/Users/<user>/.bundle/<bundle>/dev/artifacts/.internal/*.yml'
  5. To share integrations across workspace users, add top-level permissions in databricks.yml:

    YAML
    permissions:
    - group_name: 'users'
    level: CAN_VIEW

Use a registered integration​

After an integration is registered, any user in the workspace can add it to a job:

  1. Refresh the page, or navigate away and back to the Lakeflow Jobs page, to reload the list of available integrations.

  2. When adding a task to a job, click Add another task type.

  3. Find registered integrations in the Integrations section of the Add task dialog and select one.

  4. Configure the task. The form fields are generated from the integration's configuration.

note

Custom icons for integrations are not yet supported.

API reference​

This section describes the Python API for authoring integrations, the bundle task type that runs them, and the YAML schema that registers them.

Python API​

Import these objects from databricks.lakeflow.integrations to define functions and sensors.

@integration​

The @integration decorator adds metadata to a function or to a Sensor class. It generates the Lakeflow Integration YAML file that registers the integration in the workspace or user integration list.

Python
from databricks.lakeflow.integrations import integration

Function parameters and class constructor parameters become task parameters. Their string values are interpreted as JSON, so int, str, bool, dict, and list are all supported.

Sensor​

A protocol for objects that poll an external condition and either complete or defer. A Sensor is re-created between poll calls, so any state that must survive a deferral has to be persisted externally, for example in task values, workspace files, or Lakebase.

Python
from databricks.lakeflow.integrations import Context, Sensor, SensorResult

As with functions, you can add task parameters through the __init__ method.

Methods

  • poll(self, ctx: Context) -> SensorResult: Called once per attempt. Return SensorResult.completed() when the condition is met, or SensorResult.deferred(duration) to release compute and try again later.

Context​

Passed to Sensor.poll, and provides information on the task run.

Python
from databricks.lakeflow.integrations import Context

Attributes

Attribute

Type

Description

main

str

The integration's main entry point; the function or Sensor class that the task runs.

task_key

str

The key of the task running the integration.

job_run_id

int

The ID of the current job run.

task_run_id

int

The ID of the current task run.

job_id

int

The ID of the job the task belongs to.

Attribute

Type

Description

main

str

The integration's main entry point; the function or Sensor class that the task runs.

task_key

str

The key of the task running the integration.

job_run_id

int

The ID of the current job run.

task_run_id

int

The ID of the current task run.

job_id

int

The ID of the job the task belongs to.

SensorResult​

Returned by Sensor.poll to indicate whether the task has completed or should be deferred.

Python
from databricks.lakeflow.integrations import SensorResult

Field

Type

Description

status

"completed" or "deferred"

The outcome of the poll.

defer_for

datetime.timedelta

How long to defer before the next poll. Only set when deferred.

Field

Type

Description

status

"completed" or "deferred"

The outcome of the poll.

defer_for

datetime.timedelta

How long to defer before the next poll. Only set when deferred.

Methods

  • SensorResult.completed(): The condition was met and the task finishes successfully.
  • SensorResult.deferred(duration): The condition was not met. The task is rescheduled after duration, and compute is released in the meantime.

Bundle task type​

Lakeflow Integrations run as a task type called python_operator_task:

  • main: The main function, or a class that extends Sensor.
  • parameters: An array of function or class constructor parameters.

For example:

YAML
resources:
jobs:
my_function:
name: 'my_function'
tasks:
- task_key: slack
environment_key: my_environment
python_operator_task:
main: my_lakeflow_integrations.my_function
parameters:
- name: 'conn_id'
value: 'CHANGEME'
environments:
- environment_key: my_environment
spec:
environment_version: '5'
dependencies:
- ../dist/*.whl

Lakeflow Integration YAML​

The bundle template generates the Lakeflow Integration YAML automatically for any function or class annotated with @integration. The UI discovers available integrations by inspecting /Workspace/.lakeflow_integrations.yml and /Workspace/Users/<user>/.lakeflow_integrations.yml.

For example:

YAML
schema: lakeflow-integration-v0.1.0
name: Slack message
description: Post a message to a Slack channel
icon:
name: send # A known Databricks icon. See the list of available options below.
library: databricks # The only available option currently.
main: slack_operator.integrations.slack.send_message
environment:
environment_version: '5'
dependencies: # Must include the Python wheel that contains the sensor or function.
- /Workspace/Shared/integrations/slack_operator-0.1.0.whl
config:
type: object
properties:
channel: # The name of your parameter.
type: string
title: Channel # The title to render in the UI instead of the raw parameter name.
description: Channel to post to.
default: '#alerts' # Default value for new task instances.
examples: ['#my-channel'] # The first example is used as the placeholder if the field is empty.
x-ui:
widget: input # The widget to render: input, textarea, or number.
message:
type: string
title: Message
description: Message body.
examples: ['Pipeline :white_check_mark: completed']
x-ui:
widget: textarea
required: # Parameters listed here are marked as required in the UI.
- channel
- message

For the full list of icon values, expand the details below.

Available icon values

  • app
  • arrow-in
  • at
  • backup
  • bar-chart
  • beaker
  • binary
  • book
  • bookmark
  • brackets-curly
  • brackets-square
  • branch
  • briefcase
  • brush
  • bug
  • calendar
  • camera
  • catalog
  • chain
  • chart-line
  • check-circle
  • checklist
  • chip
  • clipboard
  • clock
  • cloud
  • cloud-database
  • code
  • columns
  • compass
  • connect
  • copy
  • dag
  • dashboard
  • database
  • decimal
  • dollar
  • download
  • erd
  • face-smile
  • file
  • filter
  • flag
  • flow
  • folder
  • fork
  • function
  • gear
  • gift
  • globe
  • grid
  • hash
  • history
  • home
  • image
  • ingestion
  • key
  • layer
  • leaf
  • letters
  • lightbulb
  • lightning
  • link
  • list
  • lock
  • mail
  • map
  • measure
  • megaphone
  • models
  • moon
  • notebook
  • notification
  • numbers
  • office
  • pencil
  • pie-chart
  • pipeline
  • play
  • plug
  • puzzle
  • query
  • refresh
  • robot
  • rocket
  • rows
  • school
  • search
  • send
  • share
  • shield
  • sliders
  • sparkle
  • speech-bubble
  • speedometer
  • star
  • storefront
  • stream
  • sun
  • sync
  • table
  • tag
  • target
  • terminal
  • trash
  • tree
  • trending
  • upload
  • user
  • user-group
  • visible
  • workflows
  • wrench
  • zoom-in

Billing and compute​

Integrations run on serverless compute, and billing is similar to running notebooks. Usage appears in system tables under the run name lakeflow_integrations:

SQL
SELECT *
FROM system.billing.usage
WHERE usage_metadata.job_name = 'lakeflow_integrations'

Compute is reused and released to reduce cost when a task meets any of the following conditions:

  • The deferral duration is longer than one minute.
  • Multiple parallel sensors run on behalf of the same run-as user or service principal, so the underlying compute can be shared across them.
  • The task uses serverless client version 5 or later.

When multiple job runs reuse the same compute, usage is attributed to the first job run that acquired the compute. For example, running 10 parallel sensors, each with 20 iterations and a five-minute deferral, produces about 22 records totaling roughly 0.48 DBU. The exact figure varies with the workload.

SQL
SELECT
usage_metadata.job_name,
SUM(usage_quantity) AS usage_quantity,
COUNT(*) AS records
FROM system.billing.usage
WHERE usage_metadata.job_name = 'lakeflow_integrations'
GROUP BY ALL

Frequently asked questions​

The following questions address common issues.

How do I create a Unity Catalog connection for an external HTTP service?​

Follow Connect to external HTTP services.

Can I run a compute-heavy workload?​

Resource-intensive workloads can affect other workloads that reuse the same compute, so Databricks doesn't recommend using this feature with compute-heavy workloads.