Pular para o conteúdo principal
Página não listada
Esta página não está listada. Mecanismos de busca não armazenarão nenhuma informação, e somente usuários que possuam o link direto poderão acessá-la

Integrações do Lakeflow

info

Beta

Este recurso está em Beta.

As integrações são tarefas personalizadas que você pode adicionar aos Lakeflow Jobs. Um autor escreve uma integração em Python e a registra no workspace, tornando-a disponível na caixa de diálogo Adicionar tarefa . Qualquer outro usuário pode então usá-la em um job sem escrever código.

Uma integração é uma função ou um sensor:

  • Funções realizam uma operação única, como o envio de uma notificação.
  • Sensores aguardam uma condição verificando-a em um loop. Entre as verificações, um sensor libera seu compute em vez de mantê-lo parado, para que possa aguardar de forma eficiente.

Para começar, adicione uma integração ou use uma integração registrada.

nota

Para enviar comentários ou fazer perguntas durante a pré-visualização, envie um e-mail para lakeflow-integrations-private-preview@databricks.com.

Adicionar uma integração​

Você cria integrações usando Pacotes de Automação Declarativa. Crie um pacote a partir do modelo Lakeflow Integrations , que estrutura duas integrações de exemplo que você pode modificar para criar a sua própria. Você pode criar o pacote a partir do workspace ou da CLI do Databricks.

Para criar e registrar uma integração a partir da interface do usuário do workspace:

  1. Crie um pacote a partir do padrão Lakeflow Integrations , seguindo o Tutorial: criar e implantar um pacote no workspace.

  2. Quando o bundle for criado, clique no ícone de bundle (foguete) na barra lateral esquerda e, em seguida, clique em Implantar . A implantação faz upload dos arquivos YAML de integração e wheel, e cria Jobs de exemplo.

  3. Encontre as wheels e os arquivos YAML upload em /Workspace/Users/<user>/.bundle/<bundle>/dev/artifacts/.internal.

  4. Registre as integrações upload para que apareçam na UI, usando um arquivo .lakeflow_integrations.yml:

    • Para disponibilizar uma integração a qualquer usuário, adicione-a ao /Workspace/.lakeflow_integrations.yml e conceda aos usuários do Workspace permissão para ler o arquivo.
    • Para tornar uma integração visível apenas para você, adicione-a a /Workspace/Users/<user>/.lakeflow_integrations.yml.

    Por exemplo:

    YAML
    integrations:
    - '/Workspace/Users/<user>/.bundle/<bundle>/dev/artifacts/.internal/*.yml'
  5. Para compartilhar integrações entre usuários do workspace, adicione permissões de nível superior em databricks.yml:

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

Usar uma integração registrada​

Após o registro de uma integração, qualquer usuário no workspace pode adicioná-la a um job:

  1. Refresh a página, ou saia e retorne à página do Lakeflow Jobs, para recarregar a lista de integrações disponíveis.

  2. Ao adicionar uma tarefa a um job, clique em Adicionar outro tipo de tarefa .

  3. Encontre integrações registradas na seção Integrações da caixa de diálogo Adicionar tarefa e selecione uma.

  4. Configure a tarefa. Os campos do formulário são gerados a partir da configuração da integração.

nota

Ícones personalizados para integrações ainda não são aceitos.

Referência de API​

Esta seção descreve a API Python para criação de integrações, o tipo de tarefa de pacote que as executa e o esquema YAML que as registra.

API do Python​

Importe esses objetos de databricks.lakeflow.integrations para definir funções e sensores.

@integration​

O decorador @integration adiciona metadados a uma função ou a uma classe Sensor. Ele gera o arquivo YAML de integração do Lakeflow que registra a integração no workspace ou na lista de integração do usuário.

Python
from databricks.lakeflow.integrations import integration

Os parâmetros de função e os parâmetros do construtor de classe tornam-se parâmetros da tarefa. Seus valores de string são interpretados como JSON, portanto, int, str, bool, dict e list são todos suportados.

Sensor​

Um protocolo para objetos que sondam uma condição externa e concluem ou adiam. Um Sensor é recriado entre chamadas de sondagem, portanto, qualquer estado que deva sobreviver a um adiamento precisa ser persistido externamente, por exemplo, em valores de tarefa, arquivos de workspace ou Lakebase.

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

Assim como com funções, você pode adicionar parâmetros de tarefa por meio do método __init__.

Métodos

  • poll(self, ctx: Context) -> SensorResult: Chamado uma vez por tentativa. Retorne SensorResult.completed() quando a condição for atendida, ou SensorResult.deferred(duration) para liberar o compute e tentar novamente mais tarde.

Context​

Passado para Sensor.poll e fornece informações sobre a execução da tarefa.

Python
from databricks.lakeflow.integrations import Context

Atributos

Atributo

Tipo

Descrição

main

str

O ponto de entrada principal da integração; a função ou classe Sensor que a tarefa executa.

task_key

str

A key da tarefa que executa a integração.

job_run_id

int

O ID da execução do job atual.

task_run_id

int

A ID da execução da tarefa atual.

job_id

int

O ID do job ao qual a tarefa pertence.

Atributo

Tipo

Descrição

main

str

O ponto de entrada principal da integração; a função ou classe Sensor que a tarefa executa.

task_key

str

A key da tarefa que executa a integração.

job_run_id

int

O ID da execução do job atual.

task_run_id

int

A ID da execução da tarefa atual.

job_id

int

O ID do job ao qual a tarefa pertence.

SensorResult​

Retornado por Sensor.poll para indicar se a tarefa foi concluída ou se deve ser adiada.

Python
from databricks.lakeflow.integrations import SensorResult

campo

Tipo

Descrição

status

"completed" ou "deferred"

O resultado da consulta.

defer_for

datetime.timedelta

Por quanto tempo adiar antes da próxima sondagem. Definido apenas quando adiado.

campo

Tipo

Descrição

status

"completed" ou "deferred"

O resultado da consulta.

defer_for

datetime.timedelta

Por quanto tempo adiar antes da próxima sondagem. Definido apenas quando adiado.

Métodos

  • SensorResult.completed(): a condição foi atendida e a tarefa é concluída com sucesso.
  • SensorResult.deferred(duration): a condição não foi atendida. A tarefa é reprogramada após duration, e o compute é liberado nesse meio tempo.

Tipo de tarefa de pacote​

As integrações do Lakeflow são executadas como um tipo de tarefa chamado python_operator_task:

  • main: A função principal ou uma classe que estende Sensor.
  • parameters: Uma matriz de parâmetros de função ou construtor de classe.

Por exemplo:

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

YAML de integração do Lakeflow​

O padrão de pacote gera o YAML de integração do Lakeflow automaticamente para qualquer função ou classe anotada com @integration. A UI descobre integrações disponíveis inspecionando /Workspace/.lakeflow_integrations.yml e /Workspace/Users/<user>/.lakeflow_integrations.yml.

Por exemplo:

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

Para a lista completa de valores de ícones, expanda os detalhes abaixo.

Valores de ícone disponíveis

  • 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

Cobrança e compute​

As integrações são executadas no compute serverless, e o faturamento é semelhante à execução de notebooks. O uso aparece nas tabelas do sistema sob o nome da execução lakeflow_integrations:

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

O compute é reutilizado e liberado para reduzir custos quando uma tarefa atende a qualquer uma das seguintes condições:

  • A duração do adiamento é maior que um minuto.
  • Vários sensores paralelos são executados em nome do mesmo usuário ou service principal de run-as, para que o compute subjacente possa ser compartilhado entre eles.
  • A tarefa utiliza a versão 5 ou posterior do cliente serverless.

Quando várias execuções de job reutilizam o mesmo compute, o uso é atribuído à primeira execução de job que adquiriu o compute. Por exemplo, executar 10 sensores paralelos, cada um com 20 iterações e um adiamento de cinco minutos, produz cerca de 22 registros totalizando aproximadamente 0,48 DBU. O valor exato varia de acordo com a carga de trabalho.

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

Perguntas frequentes​

As perguntas a seguir abordam problemas comuns.

Como crio uma conexão do Unity Catalog para um serviço HTTP externo?​

Siga Conectar a serviços HTTP externos.

Posso executar uma carga de trabalho com uso intenso de compute?​

Cargas de trabalho intensivas em recursos podem afetar outras cargas de trabalho que reutilizam o mesmo compute, portanto, a Databricks não recomenda o uso desse recurso com cargas de trabalho de compute pesado.

Recursos adicionais​