メインコンテンツまでスキップ
非公開のページ
このページは非公開です。 検索対象外となり、このページのリンクに直接アクセスできるユーザーのみに公開されます。

Lakeflow 統合

備考

ベータ版

この機能はベータ版です。

統合は、Lakeflow Jobsに追加できるカスタムタスクです。作成者がPythonで統合を記述し、ワークスペースに登録することで、 [タスクの追加] ダイアログから利用できるようになります。その後、他のユーザーはコードを記述することなく、ジョブ内でそれを使用できます。

インテグレーションは、関数またはセンサーのいずれかです:

  • 関数 は、通知の送信などの1回限りの操作を実行します。
  • センサー は、ループ内で条件をチェックすることによって待機します。チェックの合間、センサーはコンピュートをアイドル状態で保持するのではなく解放するため、効率的に待機できます。

起動するには、統合を追加するか、登録済みの統合を使用してください。

注記

プレビュー期間中にフィードバックを提供したり質問したりするには、lakeflow-integrations-private-preview@databricks.com までEメールでご連絡ください。

統合の追加​

宣言型オートメーションバンドルを使用して統合を作成します。 Lakeflow Integrations Template からバンドルを作成します。これには、独自の統合を構築するために変更可能な 2 つのサンプル統合が含まれています。ワークスペースまたは Databricks CLI からバンドルを作成できます。

ワークスペースUIから統合を作成および登録するには:

  1. Lakeflow Integrations Templateからバンドルを作成し、「チュートリアル: ワークスペースでバンドルを作成およびデプロイする」に従ってください。

  2. バンドルが作成されたら、左側のサイドバーにあるバンドル(ロケット)アイコンをクリックし、 [デプロイ] をクリックします。デプロイによって統合 YAML ファイルと wheel ファイルがuploadされ、ジョブの例が作成されます。

  3. アップロードされたホイールおよび YAML ファイルは /Workspace/Users/<user>/.bundle/<bundle>/dev/artifacts/.internal にあります。

  4. uploadされた統合をUIに表示させるには、.lakeflow_integrations.ymlファイルを使用して登録します:

    • 統合をすべてのユーザーが利用できるようにするには、それを /Workspace/.lakeflow_integrations.yml に追加し、ワークスペースのユーザーにファイルの読み取り権限を付与します。
    • 統合を自分のみに表示させるには、それを/Workspace/Users/<user>/.lakeflow_integrations.ymlに追加します。

    例えば:

    YAML
    integrations:
    - '/Workspace/Users/<user>/.bundle/<bundle>/dev/artifacts/.internal/*.yml'
  5. ワークスペースユーザー間で統合を共有するには、databricks.yml で最上位のアクセス許可を追加します。

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

登録済みの統合を使用する​

統合が登録されると、ワークスペース内のすべてのユーザーがそれをジョブに追加できるようになります:

  1. ページを更新するか、別のページに移動してからLakeflow Jobsページに戻り、利用可能な統合のリストを再読み込みしてください。

  2. ジョブにタスクを追加する際は、 [別のタスクタイプを追加] をクリックします。

  3. [タスクの追加] ダイアログの [インテグレーション] セクションで登録済みのインテグレーションを見つけ、選択します。

  4. タスクを構成します。フォームフィールドは、統合の構成から生成されます。

注記

統合用のカスタムアイコンはまだサポートされていません。

API リファレンス​

このセクションでは、統合を作成するためのPython API、それらを実行するバンドルタスクタイプ、およびそれらを登録するYAMLスキーマについて説明します。

Python API​

関数とセンサーを定義するために、databricks.lakeflow.integrations からこれらのオブジェクトをインポートします。

@integration​

@integrationデコレーターは、関数またはSensorクラスにメタデータを追加します。ワークスペースまたはユーザー統合リストに統合を登録するLakeflow統合YAMLファイルを生成します。

Python
from databricks.lakeflow.integrations import integration

関数パラメーターおよびクラスコンストラクターのパラメーターは、タスクパラメーターになります。それらの文字列値は JSON として解釈されるため、int、str、bool、dict、およびlistはすべてサポートされています。

Sensor​

外部条件をポーリングし、完了または延期を行うオブジェクトのためのプロトコル。Sensor はポーリング呼び出しのたびに再作成されるため、延期後も保持する必要がある状態は、タスクバリュー、ワークスペースファイル、または Lakebase などに外部で永続化する必要があります。

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

関数と同様に、__init__メソッドを使用してタスクパラメーターを追加できます。

手法

  • poll(self, ctx: Context) -> SensorResult: 試行ごとに1回呼び出されます。条件が満たされた場合は SensorResult.completed() を返し、そうでない場合は SensorResult.deferred(duration) を返してコンピュートを解放し、後でもう一度試してください。

Context​

Sensor.poll に渡され、タスク実行に関する情報を提供します。

Python
from databricks.lakeflow.integrations import Context

属性

属性

Type

説明

main

str

統合のメインエントリポイント。タスクが実行する関数または Sensor クラス。

task_key

str

統合を実行しているタスクのキー。

job_run_id

int

現在のジョブ ランのID。

task_run_id

int

現在のタスクランのID。

job_id

int

タスクが属するジョブの ID。

属性

Type

説明

main

str

統合のメインエントリポイント。タスクが実行する関数または Sensor クラス。

task_key

str

統合を実行しているタスクのキー。

job_run_id

int

現在のジョブ ランのID。

task_run_id

int

現在のタスクランのID。

job_id

int

タスクが属するジョブの ID。

SensorResult​

タスクが完了したか、または延期すべきかを示すために Sensor.poll によって返されます。

Python
from databricks.lakeflow.integrations import SensorResult

フィールド

Type

説明

status

"completed" または "deferred"

ポーリングの結果。

defer_for

datetime.timedelta

次のポーリングまでの延期時間。延期されている場合にのみ設定されます。

フィールド

Type

説明

status

"completed" または "deferred"

ポーリングの結果。

defer_for

datetime.timedelta

次のポーリングまでの延期時間。延期されている場合にのみ設定されます。

手法

  • SensorResult.completed():条件が満たされ、タスクは正常に終了しました。
  • SensorResult.deferred(duration):条件は満たされていません。タスクはduration後に再スケジュールされ、その間コンピュートは解放されます。

バンドルのタスクタイプ​

Lakeflow統合は、python_operator_taskと呼ばれるタスクタイプとして実行されます:

  • main:メイン関数、または Sensor を拡張するクラス。
  • parameters:関数またはクラスコンストラクターのパラメーターの配列。

例えば:

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 統合 YAML​

バンドルTemplateは、@integrationで注釈が付けられた関数またはクラスに対して、Lakeflow統合YAMLを自動的に生成します。UIは、/Workspace/.lakeflow_integrations.ymlおよび/Workspace/Users/<user>/.lakeflow_integrations.ymlを検査することで、利用可能な統合を発見します。

例えば:

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

アイコン値の一覧については、以下の詳細を展開してください。

使用可能なアイコン値

  • 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

請求とコンピュート​

統合はServerless コンピュート上で実行され、課金はノートブックの実行と同様です。使用状況は、システムテーブルのラン名 lakeflow_integrations の下に表示されます:

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

タスクが以下のいずれかの条件を満たす場合、コスト削減のためにコンピュートが再利用および解放されます。

  • 猶予期間は1分を超えています。
  • 複数の並列センサーが同じ実行ユーザーまたはService Principalに代わって実行されるため、ベースとなるコンピュートをそれらの間で共有できます。
  • このタスクは、サーバーレスクライアントバージョン5以降を使用します。

複数のジョブランが同じコンピュートを再利用する場合、使用量はそのコンピュートを最初に取得したジョブランに帰属します。たとえば、10 個の並列センサーを実行し、それぞれが 20 回のイテレーションと 5 分の遅延を行う場合、合計で約 0.48 DBU となる約 22 件のレコードが生成されます。正確な数値はワークロードによって異なります。

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

よくある質問​

以下の質問は、一般的な問題に対処するものです。

外部 HTTP サービス用の Unity Catalog 接続を作成するにはどうすればよいですか?​

「外部HTTPサービスへの接続」に従ってください。

コンピュート負荷の高いワークロードを実行できますか?​

リソース負荷の高いワークロードは、同じコンピュートを再利用する他のワークロードに影響を与える可能性があるため、Databricks では、コンピュート負荷の高いワークロードでこの機能を使用することをお勧めしません。

その他のリソース​