タスクの値を使用してタスク間で情報を渡す
タスク値は 、Databricksジョブ内のタスク間で任意の値を渡すことができる Databricks ユーティリティ taskValues サブユーティリティを参照します。taskValues サブユーティリティ (dbutils.jobs.taskValues)を参照してください。
1 つのタスクで dbutils.jobs.taskValues.set() を使用してキーと値のペアを指定し、タスク名とキーを使用して後続のタスクで値を参照できます。
dbutils.jobs.taskValues サブユーティリティの dbutils.jobs.taskValues.set() と dbutils.jobs.taskValues.get() は Python 関数であるため、言語として Python が選択されたノートブックでのみ使用できます。ただし、パラメーターをサポートするすべてのタスクに対して、動的値参照を使用してタスク値を参照できます。 参照タスク値を参照してください。
タスクの値を設定する
dbutils.jobs.taskValues.set()を使用して Python ノートブックのタスク値を設定します。
タスク値のキーは文字列である必要があります。 ノートブックに複数のタスク値が定義されている場合は、各キーが一意である必要があります。
タスクの値をキーに手動またはプログラムで割り当てることができます。 有効な JSON として表現できる値のみが許可されます。 値の JSON 表現のサイズは 48 KiB を超えることはできません。
たとえば、次の例では、キー の静的文字列を設定します fave_food。
dbutils.jobs.taskValues.set(key = "fave_food", value = "beans")
次の例では、ノートブックタスクパラメーターを使用して、特定の予約に関するすべての更新レコードをクエリーし、現在の予約ステータスとレコードの合計数を返します:
from pyspark.sql.functions import col
dbutils.widgets.text("booking_id", "51567", "Booking ID")
booking_id = dbutils.widgets.get("booking_id")
query = (spark.read.table("samples.wanderbricks.booking_updates")
.orderBy(col("updated_at"), ascending=False)
.where(col("booking_id") == booking_id)
.select(col("status"))
)
dbutils.jobs.taskValues.set(key = "record_count", value = query.count())
dbutils.jobs.taskValues.set(key = "booking_status", value = query.take(1)[0][0])
このパターンを使用して値のリストを渡し、For eachタスクなどの下流ロジックの調整に使用できます。「For eachタスクを使用して別のタスクをループで実行する」を参照してください。
次の例では、宛先IDの個別の値をPythonリストに抽出し、これをタスク値として設定します。
dest_list = list(spark.read.table("samples.wanderbricks.properties").select("destination_id").distinct().toPandas()["destination_id"])
dbutils.jobs.taskValues.set(key = "dest_list", value = dest_list)
参照タスクの値
Databricks では、動的値参照パターン {{tasks.<task_name>.values.<value_name>}}を使用して構成されたタスク パラメーターとしてタスク値を参照することをお勧めします。
たとえば、destination_lookupという名前のタスクのキー dest_list を持つタスク値を参照するには、構文 {{tasks.destination_lookup.values.dest_list}}を使用します。
「タスクパラメーターの構成」と「動的値の参照」を参照してください。
使う dbutils.jobs.taskValues.get
構文 dbutils.jobs.taskValues.get() では、アップストリーム タスク名を指定する必要があります。 この構文は、複数のダウンストリーム タスクでタスク値を使用できるため、タスク名が変更された場合に多数の更新が必要になるため、お勧めしません。
この構文を使用して、オプションで default 値と debugValueを指定できます。 キーが見つからない場合は、デフォルト値が使用されます。 この debugValue では、ノートブックをタスクとしてスケジュールする前に、ノートブックでの手動コード開発およびテスト中に使用する静的な値を設定できます。
次の例では、タスク名 に設定されたキー booking_status の値を取得します booking_lookup。 値 confirmed は、ノートブックを対話形式で実行している場合にのみ返されます。
booking_status = dbutils.jobs.taskValues.get(taskKey = "booking_lookup", key = "booking_status", debugValue = "confirmed")
Databricks では、キーの欠落やタスクの名前の誤りによる予期されるエラーメッセージのトラブルシューティングや防止が困難な場合があるため、デフォルト値の設定はお勧めしません。
タスクのイテレーションからの値の読み取りFor each
For eachタスク内で実行されるタスクは反復ごとに1回実行され、各タスクと同様に、各反復でdbutils.jobs.taskValues.set()を使用してタスクの値を設定できます。ダウンストリームタスクは、反復インデックスで並べ替えられた単一の返されたリストから、各反復の値読み取ることができます。
ベータ版
For eachタスクの集計された反復出力の読み取りは、ベータ版です。It is enabled by default and requires Databricks Runtime 15.4 LTS or above.
いずれかの読み取りサーフェスで、For each タスク自体ではなく、For each タスク内で実行されるネストされたタスクのキーを使用します。そのキーは、反復インデックスによって順序付けられた単一のリストとして、すべての反復からの値を返します。たとえば、process という名前のネストされたタスクの 3 回の反復がそれぞれ result を設定する場合:
results = dbutils.jobs.taskValues.get(taskKey = "process", key = "result")
# results == [1, 2, 3]
動的値参照 {{tasks.process.values.result}} は、同じ配列に解決されます。
キーを設定しなかったイテレーションは、結果内でその位置を維持します。dbutils.jobs.taskValues.get からは None (JSON null) として表示され、動的値参照では null として表示されます。たとえば、3 回のイテレーションの 2 回目で result が設定されなかった場合、results は [1, None, 3] になります。
集約された参照には、単一のパラメーター値で組み立てられた配列が約 48 KB (49,344 文字) を超えてはならないことと、単一のパラメーター値に含めることができる集約されたタスク値参照は最大 3 つであることの 2 つの制限があります。
タスクの値を表示する
各実行におけるタスク値の戻り値は、**タスク実行の詳細**の**出力**ペインに表示されます。「タスク実行履歴の表示」を参照してください。