さらに @tasks を使う
ワークフローの実装を続けるために、ダウンロードしたファイルの変更内容を解析・保存するステップを追加します。DAG process_sales はすでに定義されており、pull_file タスクも追加済みです。今回は、Python 関数 parse_file(inputfile, outputfile) があらかじめ定義されています。
Airflow のタスクを実装する際、個々のステップの内容を完全に理解していなくても問題ありません。各ステップを Airflow の構造の中に正しく組み込む方法を理解していれば、目的のワークフローを実装できます。
この演習はコースの一部です
Python で学ぶ Apache Airflow 入門
演習の手順
parse_fileメソッドを使って Airflow タスクを作成してください。- 必要な引数を指定してタスクを呼び出してください。
実践的なインタラクティブ演習
このサンプルコードを完成させて、この演習に挑戦してみましょう。
@dag(dag_id='process_sales')
def process_sales():
# Decorate parse_file as a task
____
def parse_file(inputfile: str, outputfile: str):
with open(inputfile) as infile:
data = json.load(infile)
with open(outputfile, 'w') as outfile:
json.dump(data, outfile)
pull_file('http://dataserver/sales.json', 'latestsales.json')
# Call the parse_file task
____('latestsales.json', 'latestsales_parsed.json')
process_sales()