เพิ่ม @tasks
เพื่อพัฒนา workflow ต่อไป คุณต้องเพิ่มขั้นตอนสำหรับแปลงและบันทึกการเปลี่ยนแปลงของไฟล์ที่ดาวน์โหลดมา Dag process_sales ถูกกำหนดไว้แล้วและมี task pull_file อยู่ในนั้นแล้ว ในกรณีนี้ Python function ถูกเตรียมไว้ให้แล้วในชื่อ parse_file(inputfile, outputfile)
โปรดทราบว่าในการ implement Airflow tasks จริง คุณอาจไม่จำเป็นต้องเข้าใจรายละเอียดของแต่ละขั้นตอนทั้งหมด ขอเพียงเข้าใจวิธีนำขั้นตอนเหล่านั้นมาใส่ใน Airflow structure ก็เพียงพอที่จะสร้าง workflow ตามต้องการได้
แบบฝึกหัดนี้เป็นส่วนหนึ่งของหลักสูตร
Apache Airflow เบื้องต้นด้วย Python
คำแนะนำการฝึกหัด
- สร้าง Airflow task โดยใช้ method
parse_file - เรียกใช้ task พร้อมอาร์กิวเมนต์ที่จำเป็น
แบบฝึกหัดเชิงโต้ตอบแบบลงมือทำ
ลองทำแบบฝึกหัดนี้โดยเติมโค้ดตัวอย่างนี้ให้สมบูรณ์
@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()