เริ่มต้นใช้งานเริ่มต้นใช้งานได้ฟรี

เพิ่ม @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()
แก้ไขและรันโค้ด