Lakeflow 統合
ベータ版
この機能はベータ版です。
統合は、Lakeflow Jobsに追加できるカスタムタスクです。作成者がPythonで統合を記述し、ワークスペースに登録することで、 [タスクの追加] ダイアログから利用できるようになります。その後、他のユーザーはコードを記述することなく、ジョブ内でそれを使用できます。
インテグレーションは、関数またはセンサーのいずれかです:
- 関数 は、通知の送信などの1回限りの操作を実行します。
- センサー は、ループ内で条件をチェックすることによって待機します。チェックの合間、センサーはコンピュートをアイドル状態で保持するのではなく解放するため、効率的に待機できます。
起動するには、統合を追加するか、登録済みの統合を使用してください。
プレビュー期間中にフィードバックを提供したり質問したりするには、lakeflow-integrations-private-preview@databricks.com までEメールでご連絡ください。
統合の追加
宣言型オートメーションバンドルを使用して統合を作成します。 Lakeflow Integrations Template からバンドルを作成します。これには、独自の統合を構築するために変更可能な 2 つのサンプル統合が含まれています。ワークスペースまたは Databricks CLI からバンドルを作成できます。
- Workspace
- Databricks CLI
ワークスペースUIから統合を作成および登録するには:
-
Lakeflow Integrations Templateからバンドルを作成し、「チュートリアル: ワークスペースでバンドルを作成およびデプロイする」に従ってください。
-
バンドルが作成されたら、左側のサイドバーにあるバンドル(ロケット)アイコンをクリックし、 [デプロイ] をクリックします。デプロイによって統合 YAML ファイルと wheel ファイルがuploadされ、ジョブの例が作成されます。
-
アップロードされたホイールおよび YAML ファイルは
/Workspace/Users/<user>/.bundle/<bundle>/dev/artifacts/.internalにあります。 -
uploadされた統合をUIに表示させるには、
.lakeflow_integrations.ymlファイルを使用して登録します:- 統合をすべてのユーザーが利用できるようにするには、それを
/Workspace/.lakeflow_integrations.ymlに追加し、ワークスペースのユーザーにファイルの読み取り権限を付与します。 - 統合を自分のみに表示させるには、それを
/Workspace/Users/<user>/.lakeflow_integrations.ymlに追加します。
例えば:
YAMLintegrations:
- '/Workspace/Users/<user>/.bundle/<bundle>/dev/artifacts/.internal/*.yml' - 統合をすべてのユーザーが利用できるようにするには、それを
-
ワークスペースユーザー間で統合を共有するには、
databricks.ymlで最上位のアクセス許可を追加します。YAMLpermissions:
- group_name: 'users'
level: CAN_VIEW
Databricks CLI から統合を作成および登録するには:
-
Databricks CLIバージョン1.7.0以降をインストールします。Databricks CLI のインストールまたは更新を参照してください。
Bashdatabricks version
# Databricks CLI v1.7.0 -
まだ認証していない場合は、ワークスペースを認証してください:
Bashdatabricks configure -
uv をインストールします。uv のインストールを参照してください。バンドルはパッケージングに uv を使用します。
-
Lakeflow Integrations Templateからバンドルを作成します:
Bashdatabricks bundle init lakeflow-integrations -
プロジェクトディレクトリに移動し、バンドルをデプロイします:
Bashcd my_lakeflow_integrations
databricks bundle deployデプロイメントは統合 YAML と Python wheel をビルドし、それらをワークスペースにアップロードします。
-
統合YAMLファイルとwheelファイルがuploadされたことを確認します:
Bashdatabricks workspace list /Workspace/Users/<user>/.bundle/my_lakeflow_integrations/dev/artifacts/.internal -
uploadされた統合をUIに表示させるには、
.lakeflow_integrations.ymlファイルを使用して登録します:- 統合をすべてのユーザーが利用できるようにするには、それを
/Workspace/.lakeflow_integrations.ymlに追加し、ワークスペースのユーザーにファイルの読み取り権限を付与します。 - 統合を自分のみに表示させるには、それを
/Workspace/Users/<user>/.lakeflow_integrations.ymlに追加します。
例えば:
YAMLintegrations:
- '/Workspace/Users/<user>/.bundle/my_lakeflow_integrations/dev/artifacts/.internal/*.yml'Databricks CLI を使用してファイルをuploadします:
Bashdatabricks workspace import /Workspace/.lakeflow_integrations.yml --file .lakeflow_integrations.yml --format AUTO - 統合をすべてのユーザーが利用できるようにするには、それを
-
ワークスペースユーザー間で統合を共有するには、
databricks.ymlで最上位のアクセス許可を追加します。YAMLpermissions:
- group_name: 'users'
level: CAN_VIEW -
このTemplateでは、確認や実行が可能なジョブの例も作成されます。
Bashdatabricks bundle summary
databricks bundle run
登録済みの統合を使用する
統合が登録されると、ワークスペース内のすべてのユーザーがそれをジョブに追加できるようになります:
-
ページを更新するか、別のページに移動してからLakeflow Jobsページに戻り、利用可能な統合のリストを再読み込みしてください。
-
ジョブにタスクを追加する際は、 [別のタスクタイプを追加] をクリックします。
-
[タスクの追加] ダイアログの [インテグレーション] セクションで登録済みのインテグレーションを見つけ、選択します。
-
タスクを構成します。フォームフィールドは、統合の構成から生成されます。
統合用のカスタムアイコンはまだサポートされていません。
API リファレンス
このセクションでは、統合を作成するためのPython API、それらを実行するバンドルタスクタイプ、およびそれらを登録するYAMLスキーマについて説明します。
Python API
関数とセンサーを定義するために、databricks.lakeflow.integrations からこれらのオブジェクトをインポートします。
@integration
@integrationデコレーターは、関数またはSensorクラスにメタデータを追加します。ワークスペースまたはユーザー統合リストに統合を登録するLakeflow統合YAMLファイルを生成します。
from databricks.lakeflow.integrations import integration
関数パラメーターおよびクラスコンストラクターのパラメーターは、タスクパラメーターになります。それらの文字列値は JSON として解釈されるため、int、str、bool、dict、およびlistはすべてサポートされています。
Sensor
外部条件をポーリングし、完了または延期を行うオブジェクトのためのプロトコル。Sensor はポーリング呼び出しのたびに再作成されるため、延期後も保持する必要がある状態は、タスクバリュー、ワークスペースファイル、または Lakebase などに外部で永続化する必要があります。
from databricks.lakeflow.integrations import Context, Sensor, SensorResult
関数と同様に、__init__メソッドを使用してタスクパラメーターを追加できます。
手法
poll(self, ctx: Context) -> SensorResult: 試行ごとに1回呼び出されます。条件が満たされた場合はSensorResult.completed()を返し、そうでない場合はSensorResult.deferred(duration)を返してコンピュートを解放し、後でもう一度試してください。
Context
Sensor.poll に渡され、タスク実行に関する情報を提供します。
from databricks.lakeflow.integrations import Context
属性
属性 | Type | 説明 |
|---|---|---|
|
| 統合のメインエントリポイント。タスクが実行する関数または |
|
| 統合を実行しているタスクのキー。 |
|
| 現在のジョブ ランのID。 |
|
| 現在のタスクランのID。 |
|
| タスクが属するジョブの ID。 |
SensorResult
タスクが完了したか、または延期すべきかを示すために Sensor.poll によって返されます。
from databricks.lakeflow.integrations import SensorResult
フィールド | Type | 説明 |
|---|---|---|
|
| ポーリングの結果。 |
|
| 次のポーリングまでの延期時間。延期されている場合にのみ設定されます。 |
手法
SensorResult.completed():条件が満たされ、タスクは正常に終了しました。SensorResult.deferred(duration):条件は満たされていません。タスクはduration後に再スケジュールされ、その間コンピュートは解放されます。
バンドルのタスクタイプ
Lakeflow統合は、python_operator_taskと呼ばれるタスクタイプとして実行されます:
main:メイン関数、またはSensorを拡張するクラス。parameters:関数またはクラスコンストラクターのパラメーターの配列。
例えば:
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を検査することで、利用可能な統合を発見します。
例えば:
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
アイコン値の一覧については、以下の詳細を展開してください。
使用可能なアイコン値
apparrow-inatbackupbar-chartbeakerbinarybookbookmarkbrackets-curlybrackets-squarebranchbriefcasebrushbugcalendarcameracatalogchainchart-linecheck-circlechecklistchipclipboardclockcloudcloud-databasecodecolumnscompassconnectcopydagdashboarddatabasedecimaldollardownloaderdface-smilefilefilterflagflowfolderforkfunctiongeargiftglobegridhashhistoryhomeimageingestionkeylayerleafletterslightbulblightninglinklistlockmailmapmeasuremegaphonemodelsmoonnotebooknotificationnumbersofficepencilpie-chartpipelineplayplugpuzzlequeryrefreshrobotrocketrowsschoolsearchsendshareshieldsliderssparklespeech-bubblespeedometerstarstorefrontstreamsunsynctabletagtargetterminaltrashtreetrendinguploaduseruser-groupvisibleworkflowswrenchzoom-in
請求とコンピュート
統合はServerless コンピュート上で実行され、課金はノートブックの実行と同様です。使用状況は、システムテーブルのラン名 lakeflow_integrations の下に表示されます:
SELECT *
FROM system.billing.usage
WHERE usage_metadata.job_name = 'lakeflow_integrations'
タスクが以下のいずれかの条件を満たす場合、コスト削減のためにコンピュートが再利用および解放されます。
- 猶予期間は1分を超えています。
- 複数の並列センサーが同じ実行ユーザーまたはService Principalに代わって実行されるため、ベースとなるコンピュートをそれらの間で共有できます。
- このタスクは、サーバーレスクライアントバージョン5以降を使用します。
複数のジョブランが同じコンピュートを再利用する場合、使用量はそのコンピュートを最初に取得したジョブランに帰属します。たとえば、10 個の並列センサーを実行し、それぞれが 20 回のイテレーションと 5 分の遅延を行う場合、合計で約 0.48 DBU となる約 22 件のレコードが生成されます。正確な数値はワークロードによって異なります。
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 では、コンピュート負荷の高いワークロードでこの機能を使用することをお勧めしません。