Chuyển đến nội dung chính

5 lựa chọn thay thế Airflow hàng đầu cho điều phối dữ liệu (Kèm ví dụ mã)

Khám phá năm lựa chọn thay thế Airflow để điều phối dữ liệu với ví dụ mã xây dựng, chạy và trực quan hóa một pipeline ETL đơn giản.
Đã cập nhật 31 thg 8, 2026  · 13 phút đọc

Khám phá với AI

ChatGPTClaudePerplexity

Choose an Airflow Alternatives meme template

Hình ảnh do tác giả cung cấp.

Apache Airflow là một công cụ điều phối dữ liệu mã nguồn mở phổ biến, được thiết kế để xây dựng, lập lịch và giám sát các pipeline dữ liệu. Công cụ này có bảng điều khiển giúp quản lý trạng thái của quy trình làm việc, khiến nó trở thành lựa chọn hoàn hảo cho hầu hết nhu cầu workflow.

Tuy nhiên, Airflow thiếu một số tính năng quan trọng có thể mang tính sống còn đối với các yêu cầu điều phối dữ liệu phức tạp, hiện đại.

Trong hướng dẫn này, chúng ta sẽ khám phá năm lựa chọn thay thế Airflow với các khả năng nâng cao và khắc phục một số hạn chế của nó. Bên cạnh đó, chúng ta sẽ học cách xây dựng một pipeline ETL đơn giản bằng từng công cụ, chạy nó và trực quan hóa trên bảng điều khiển của chúng.

Vì sao nên chọn một giải pháp thay thế Airflow? 

Airflow là công cụ mạnh mẽ cho nhiều workflow dữ liệu, nhưng có một số hạn chế có thể khiến các công ty cân nhắc lựa chọn khác. 

Dưới đây là một số lý do bạn có thể chọn giải pháp thay thế:

  1. Đường cong học tập dốc: Airflow có thể khó học, đặc biệt với người mới dùng công cụ quản lý workflow.
  2. Bảo trì: Cần bảo trì đáng kể, nhất là ở các triển khai quy mô lớn.
  3. Tài liệu chưa đầy đủ: Nhiều người dùng phản ánh vấn đề về tài liệu, gây khó khăn khi xử lý sự cố hoặc tìm hiểu tính năng mới. 
  4. Tốn tài nguyên: Airflow có thể ngốn tài nguyên, đòi hỏi nhiều tính toán và bộ nhớ để chạy hiệu quả.
  5. Hạn chế linh hoạt cho người không dùng Python: Triết lý workflow dưới dạng code phụ thuộc nặng vào Python, có thể loại trừ các chuyên gia nghiệp vụ không thành thạo lập trình.
  6. Khả năng mở rộng: Một số người dùng gặp khó khăn khi mở rộng Airflow cho các workflow lớn.
  7. Xử lý thời gian thực hạn chế: Airflow chủ yếu được thiết kế cho xử lý theo lô, không phải luồng dữ liệu thời gian thực.

Trước khi đi vào phần viết mã với các công cụ điều phối dữ liệu khác, điều quan trọng là học cách viết pipeline dữ liệu bằng Apache Airflow qua hướng dẫn Bắt đầu với Apache Airflow, để bạn có thể so sánh công bằng các lựa chọn thay thế.

Nếu bạn hoàn toàn mới với Airflow, hãy cân nhắc khóa học ngắn Giới thiệu về Airflow trong Python để học những kiến thức cơ bản về xây dựng và lập lịch pipeline dữ liệu.

5 lựa chọn thay thế Airflow tốt nhất cho điều phối dữ liệu

Giờ hãy mô tả 5 lựa chọn thay thế hàng đầu cho Airflow và minh họa cách sử dụng chúng với các ví dụ mã thực tiễn.

1. Prefect

Prefect là công cụ điều phối workflow bằng Python mã nguồn mở dành cho kỹ sư dữ liệu và máy học hiện đại. Công cụ cung cấp API đơn giản cho phép bạn nhanh chóng xây dựng pipeline dữ liệu và quản lý qua một bảng điều khiển tương tác. 

Perfect cung cấp mô hình thực thi lai, nghĩa là bạn có thể triển khai workflow lên đám mây và chạy ở đó hoặc dùng kho địa phương.

So với Airflow, Prefect có các tính năng nâng cao như phụ thuộc tác vụ tự động, trigger dựa trên sự kiện, thông báo tích hợp, hạ tầng chuyên biệt cho workflow và chia sẻ dữ liệu giữa các tác vụ. Những khả năng này giúp quản lý workflow phức tạp một cách hiệu quả.

Prefect đơn giản nhưng đi kèm các tính năng mạnh. Mình chỉ mất 5 phút để chạy ví dụ. Mình đặc biệt thích thiết kế UI của dashboard, cách thiết lập thông báo, chạy lại pipeline, quản lý và giám sát mọi thứ qua Dashboard.

Abid Ali AwanAuthor

Đọc bài blog Airflow vs Prefect: Lựa chọn nào phù hợp cho workflow dữ liệu của bạn để tìm hiểu so sánh chi tiết giữa hai công cụ điều phối dữ liệu này. 

Bắt đầu với Prefect

Chúng ta sẽ bắt đầu dự án Prefect bằng cách cài đặt gói Python. Chạy lệnh sau trong terminal.

$ pip install -U prefect

Sau đó, tạo một script Python tên prefect_etl.py và viết đoạn mã sau.

from prefect import task, flow
import pandas as pd

# Extract data
@task
def extract_data():
    # Simulating data extraction
    data = {
        "name": ["Alice", "Bob", "Charlie"],
        "age": [25, 30, 35],
        "city": ["New York", "Los Angeles", "Chicago"]
    }
    df = pd.DataFrame(data)
    return df

# Transform data
@task
def transform_data(df: pd.DataFrame):
    # Example transformation: adding a new column
    df["age_plus_ten"] = df["age"] + 10
    return df

# Load data
@task
def load_data(df: pd.DataFrame):
    # Simulating data load
    print("Loading data to target destination:")
    print(df)

# Defining the flow
@flow(log_prints=True)
def etl():
    raw_data = extract_data()
    transformed_data = transform_data(raw_data)
    load_data(transformed_data)

# Running the flow
if __name__ == "__main__":
    etl()

Đoạn mã trên định nghĩa các hàm tác vụ extract_data(), transform_data(),load_data() rồi thực thi tuần tự trong một hàm flow tên etl(). Các hàm này được tạo bằng các decorator của Prefect trong Python. 

Tóm lại, chúng ta tạo một pandas DataFrame, biến đổi nó, rồi hiển thị kết quả cuối bằng print. Đây là cách đơn giản để mô phỏng một pipeline ETL.

Để thực thi workflow, chỉ cần chạy script Python bằng lệnh sau.

$ python prefect_etl.py 

Như ta thấy, lần chạy workflow đã hoàn tất thành công.

Prefect flow run logs

Logs lần chạy flow của Prefect.

Triển khai flow

Giờ ta sẽ triển khai workflow để có thể chạy theo lịch hoặc kích hoạt theo sự kiện. Việc triển khai flow cũng cho phép giám sát và quản lý nhiều workflow tập trung.

Để triển khai flow, chúng ta dùng Prefect CLI. Hàm deploy yêu cầu tên file Python, tên hàm flow trong file, và tên deployment. Trong trường hợp này, chúng ta gọi deployment là “simple_etl”.

$ prefect deploy prefect_etl.py:etl -n 'simple_etl'

Sau khi chạy script trên trong terminal, có thể bạn sẽ nhận thông báo chưa có worker pool để chạy deployment. Để tạo worker pool, dùng lệnh sau.

$ prefect worker start --pool 'datacamp'

Khi đã có worker pool, hãy mở một cửa sổ terminal khác và chạy deployment. Lệnh prefect deployment run yêu cầu đối số “<tên-hàm-flow>/<tên-deployment>” như bên dưới.

$ prefect deployment run 'etl/simple_etl

Sau khi chạy deployment, bạn sẽ nhận được thông báo rằng workflow đang chạy. Thông thường, lần chạy flow được tạo sẽ được gán tên ngẫu nhiên, trong trường hợp của tôi là witty-lorikeet.

Creating flow run for deployment 'etl/simple_etl'...
Created flow run 'witty-lorikeet'.
└── UUID: 4e0495b0-9c7e-4ed8-b9ab-5160994dc7f0
└── Parameters: {}
└── Job Variables: {}
└── Scheduled start time: 2024-06-22 14:05:01 PKT (now)
└── URL: <no dashboard available>

Để xem log đầy đủ, chuyển lại cửa sổ terminal nơi bạn đã khởi động worker pool.

Prefect flow run summary

Tóm tắt lần chạy flow của Prefect.

Bạn cần khởi động web server của Prefect để trực quan hóa lần chạy flow thân thiện hơn và quản lý các workflow khác.

$ prefect server start 

Sau khi thực thi lệnh trên, bạn sẽ được chuyển hướng đến bảng điều khiển Prefect. Hoặc có thể truy cập trực tiếp địa chỉ http://127.0.0.1:4200 trong trình duyệt.

Prefect web server UI

Giao diện web server của Prefect

Bảng điều khiển cho phép bạn chạy lại workflow, xem log, kiểm tra work pool, thiết lập thông báo và chọn các tùy chọn nâng cao khác. Đây là giải pháp hoàn chỉnh cho nhu cầu điều phối dữ liệu hiện đại.

Để học cách xây dựng và thực thi các pipeline máy học bằng Prefect, bạn có thể theo dõi hướng dẫn Sử dụng Prefect cho workflow máy học.

2. Dagster

Dasgter là một framework mã nguồn mở dành cho kỹ sư dữ liệu để định nghĩa, lập lịch và giám sát các pipeline dữ liệu. Công cụ có khả năng mở rộng cao và tạo điều kiện hợp tác giữa nhiều nhóm dữ liệu. 

Dagster cho phép người dùng định nghĩa các tài sản dữ liệu dưới dạng hàm Python sử dụng decorator. Khi các tài sản được định nghĩa, người dùng có thể thực thi mượt mà thông qua lập lịch hoặc trigger dựa trên sự kiện.

So với Airflow, Dagster cho phép phát triển, kiểm thử và rà soát pipeline cục bộ, cung cấp cách tiếp cận điều phối dựa trên tài sản, và thuần đám mây/containers.

Thay vì nghĩ về workflow theo bước và luồng, tôi phải thay đổi cách nghĩ và xây dựng pipeline bằng các tài sản dữ liệu. Ngoài điều đó ra, việc xây và chạy một pipeline ETL đơn giản khá dễ. Webserver tương đối tối giản nhưng cung cấp đủ thông tin để giám sát tài sản, lần chạy và deployment.

Abid Ali AwanAuthor

Bắt đầu với Dagster

Chúng ta sẽ tạo một pipeline ETL đơn giản, thực thi và trực quan hóa nó bằng web server của Dagster. Tương tự bảng điều khiển của Prefect, web server Dagster cung cấp cách thức tập trung để giám sát nhiều workflow, lập lịch các lần chạy và tài sản.

Bắt đầu bằng cách cài đặt gói Python.

$ pip install dagster -q

Sau đó, tạo ba hàm Python để trích xuất, biến đổi và nạp dữ liệu. Các hàm này được đặt tên create_dirty_data()clean_data(), load_cleaned_data() trong mã. Sử dụng decorator @asset, chúng ta sẽ khai báo các hàm là tài sản dữ liệu trong Dagster.

Tiếp theo, chúng ta tạo job tài sản (biến job) sử dụng tất cả tài sản (biến all_assets) rồi tạo asset definition (biến defs). 

Bạn có thể bỏ qua phần asset definition, nhưng nó quan trọng nếu muốn lập lịch lần chạy, chạy nhiều job, và thiết lập sensor.

import pandas as pd
import numpy as np
from dagster import asset, Definitions, define_asset_job, materialize

@asset
def create_dirty_data():
    # Create a sample DataFrame with dirty data
    data = {
        'Name': [' John Doe ', 'Jane Smith', 'Bob Johnson ', '  Alice Brown'],
        'Age': [30, np.nan, 40, 35],
        'City': ['New York', 'los angeles', 'CHICAGO', 'Houston'],
        'Salary': ['50,000', '60000', '75,000', 'invalid']
    }
    df = pd.DataFrame(data)
    
    # Save the DataFrame to a CSV file
    dirty_file_path = 'dag_data/dirty_data.csv'
    df.to_csv(dirty_file_path, index=False)
    
    return dirty_file_path

@asset
def clean_data(create_dirty_data):
    # Read the dirty CSV file
    df = pd.read_csv(create_dirty_data)
    
    # Clean the data
    df['Name'] = df['Name'].str.strip()
    df['Age'] = pd.to_numeric(df['Age'], errors='coerce').fillna(df['Age'].mean())
    df['City'] = df['City'].str.upper()
    df['Salary'] = df['Salary'].replace('[\$,]', '', regex=True)
    df['Salary'] = pd.to_numeric(df['Salary'], errors='coerce').fillna(0)
    
    # Calculate average salary
    avg_salary = df['Salary'].mean()
    
    # Save the cleaned DataFrame to a new CSV file
    cleaned_file_path = 'dag_data/cleaned_data.csv'
    df.to_csv(cleaned_file_path, index=False)
    
    return {
        'cleaned_file_path': cleaned_file_path,
        'avg_salary': avg_salary
    }

@asset
def load_cleaned_data(clean_data):
    cleaned_file_path = clean_data['cleaned_file_path']
    avg_salary = clean_data['avg_salary']
    
    # Read the cleaned CSV file to verify
    df = pd.read_csv(cleaned_file_path)
    
    print({
        'num_rows': len(df),
        'num_columns': len(df.columns),
        'avg_salary': avg_salary
    })

# Define all assets
all_assets = [create_dirty_data, clean_data, load_cleaned_data]

# Create a job that will materialize all assets
job = define_asset_job("all_assets_job", selection=all_assets)

# Create Definitions object
defs = Definitions(
    assets=all_assets,
    jobs=[job]
)

if __name__ == "__main__":
    result = materialize(all_assets)
    print("Pipeline execution result:", result.success)

Bạn có thể chạy đoạn mã trên trong Jupyter Notebook hoặc tạo file Python rồi chạy. 

Sau khi thực thi mã, chúng ta sẽ nhận log đầy đủ của lần chạy workflow. 

Dagster execution summary

Web server Dagster

Để trực quan hóa các tài sản và lần chạy job, chúng ta cần cài đặt và chạy web server của Dagster. Web server cho phép chạy job, materialize từng tài sản và giám sát nhiều job cùng lúc.

$ pip install dagster-webserver

Để khởi động Dagster server, chúng ta dùng Daster CLI và cung cấp vị trí file Python. Ở đây, tôi đặt tên file là dagster_pipe.py.

$ dagster dev -f dagster_pipe.py  

Lệnh trên sẽ tự động mở web server trên trình duyệt của bạn. Hoặc bạn có thể truy cập trực tiếp địa chỉ http://127.0.0.1:3000 trong trình duyệt.

Dagster Web server

Giao diện web server Dagster.

Đến giờ chúng ta mới chỉ triển khai job. Để chạy workflow, vào tab “Runs” và nhấp nút “Launch a new run”. 

Lần chạy sẽ hoàn tất thành công! Để xem log, nhấp vào ID của lần chạy bạn quan tâm.

Dagster runs detailed view

Log lần chạy của Dagster.

3. Mage AI

Mage AI là framework điều phối dữ liệu lai mã nguồn mở. “Lai” nghĩa là bạn có được sự linh hoạt của Jupyter Notebook và khả năng kiểm soát của mã mô-đun. 

Bất kỳ ai, kể cả người chỉ biết Python ở mức hạn chế, cũng có thể xây dựng, chạy và giám sát pipeline dữ liệu. Thay vì viết và chạy một file Python trực tiếp, bạn sẽ tạo dự án Mage AI và mở nó trong dashboard, nơi bạn có thể xây, chạy và quản lý các pipeline dữ liệu.

So với Airflow, Mage AI cung cấp giao diện thân thiện và dễ dùng, là lựa chọn tuyệt vời cho người mới với kỹ thuật dữ liệu. Nó được thiết kế với khả năng mở rộng và có thể xử lý lượng dữ liệu lớn cùng cấu trúc pipeline phức tạp một cách hiệu quả.

Mình thấy khá lạ vì nó hoàn toàn khác những gì mình quen thuộc. Mình phải cài đặt và mở UI web của Mage AI. Đáng lẽ phải đơn giản, nhưng mình thấy khó để xây và chạy pipeline ETL. Mặt khác, mình hiểu vì sao thiết kế độc đáo này có thể hấp dẫn người mới vào nghề, vì về cơ bản chỉ là kéo thả và bấm nút.

Abid Ali AwanAuthor

Bắt đầu với Mage AI

Bắt đầu với Mage AI khá đơn giản. Chúng ta chỉ cần cài gói Python của Mage AI.

$ pip install mage-ai

Và khởi chạy dự án Mage AI. 

$ mage start mage_ai_etl 

Lệnh trên sẽ khởi động web server. Như đã đề cập, mọi việc chỉnh sửa code, chạy job và giám sát job đều thực hiện qua UI của Mage AI.

Mage AI UI

Giao diện Mage AI.

Nhấp “+ New pipeline” để tạo pipeline ETL đầu tiên. Tôi đặt tên là “simple_etl”.

Creating the new pipeline in Mage AI

Tạo pipeline mới trong Mage AI.

Sau đó, giao diện sẽ yêu cầu bạn thêm một mô-đun để bắt đầu viết mã. Chọn mô-đun “Data Loader” và viết đoạn Python sau. 

Ở đây, chúng ta khai báo hàm create_sample_csv(), là bước đầu trong pipeline. Chúng ta dùng decorator @data_loader của Mage AI. Chúng ta cũng định nghĩa hàm test_output() để kiểm tra xem đầu ra có tồn tại hay không. Điều này giúp quản lý phụ thuộc tác vụ.

import io
import pandas as pd

if 'data_loader' not in globals():
    from mage_ai.data_preparation.decorators import data_loader
if 'test' not in globals():
    from mage_ai.data_preparation.decorators import test

@data_loader
def create_sample_csv() -> pd.DataFrame:
    """
    Create a sample CSV file with duplicates and missing values
    """
    csv_data = """
category,product,quantity,price
Electronics,Laptop,5,1000
Electronics,Smartphone,10,500
Clothing,T-shirt,50,20
Clothing,Jeans,30,50
Books,Novel,100,15
Books,Textbook,20,80
Electronics,Laptop,5,1000
Clothing,T-shirt,,20
Electronics,Tablet,,300
Books,Magazine,25,
"""
    return pd.read_csv(io.StringIO(csv_data.strip()))
   
@test
def test_output(df) -> None:
    """
    Template code for testing the output of the block.
    """
    assert df is not None, 'The output is undefined'

Creating the data loader block in Mage AI

Tạo khối data loader trong Mage AI.

Tiếp theo, tạo mô-đun “Transformer” và thêm hàm clean_data() như trong mã bên dưới. 

Bạn có thể bỏ qua hàm test(); bạn chỉ cần thêm hàm transformer chính, clean_data().

import pandas as pd

if 'transformer' not in globals():
    from mage_ai.data_preparation.decorators import transformer
if 'test' not in globals():
    from mage_ai.data_preparation.decorators import test

@transformer
def clean_data(df: pd.DataFrame) -> pd.DataFrame:
    """
    Clean and transform the data
    """
    # Remove duplicates
    df = df.drop_duplicates()
    # Fill missing values with 0
    df = df.fillna(0)
    return df

@test
def test_output(df) -> None:
    """
    Template code for testing the output of the block.
    """
    assert df is not None, 'The output is undefined'

Tương tự, tạo mô-đun “Data Exporter” và thêm đoạn mã sau. Mã khai báo hàm nạp dữ liệu export_data_to_csv(), lưu dữ liệu đã biến đổi vào tệp CSV. 

import pandas as pd

if 'data_exporter' not in globals():
    from mage_ai.data_preparation.decorators import data_exporter

@data_exporter
def export_data_to_csv(df: pd.DataFrame) -> None:
    """
    Export the processed data to a CSV file
    """
    df.to_csv('output_data.csv', index=False)
    print("Data exported successfully to output_data.csv")

Để chạy pipeline, vào tab “Trigger”, và nhấp “Run@once”.

Running the pipeline in Mage AI

Chạy pipeline trong Mage AI.

Để xem log của lần chạy, vào tab “Runs” và nhấp nút “Logs” của pipeline vừa chạy.

Mage AI flow run logs

Logs lần chạy flow của Mage AI.

4. Kedro

Kedro là một framework điều phối dữ liệu mã nguồn mở phổ biến khác, hơi khác so với các công cụ còn lại. Nó được tạo ra cho kỹ sư máy học và vay mượn nhiều khái niệm từ kỹ nghệ phần mềm để áp dụng vào dự án máy học.

Kedro được thiết kế có tính mô-đun cao, nghĩa là ngay cả để xuất một tập dữ liệu, bạn cũng phải tạo data catalog chỉ định vị trí và loại dữ liệu, đảm bảo quản lý dữ liệu chuẩn hóa và hiệu quả xuyên suốt pipeline.

Để hiểu Kedro phù hợp thế nào trong hệ sinh thái máy học, bạn có thể khám phá các công cụ MLOps khác bằng cách đọc bài 25 công cụ MLOps hàng đầu bạn cần biết trong 2024.

So với Airflow, API của Kedro đơn giản hơn để xây pipeline dữ liệu. Công cụ tập trung nhiều hơn vào kỹ nghệ máy học và cung cấp phân loại cùng versioning dữ liệu.

Phần viết mã khá trực diện, nhưng phát sinh vấn đề khi bạn muốn thực thi pipeline. Bạn phải tạo data catalog, đăng ký pipeline và nắm được cấu trúc dự án Kedro. Tôi cho rằng nó khó hơn so với Dagster và Prefect. Tuy vậy, tôi hiểu vì sao thiết kế như vậy: để khiến pipeline dữ liệu của bạn tin cậy và ít lỗi.

Abid Ali AwanAuthor

Bắt đầu với Kedro

Xây dựng pipeline dữ liệu với Kedro là một “trò chơi” khác. Framework có tính mô-đun và bạn cần hiểu cấu trúc dự án cùng các bước liên quan để thực thi workflow thành công. 

Bắt đầu bằng cách cài đặt gói Python của Kedro. 

$ pip install kedro

Khởi tạo dự án Kedro. 

$ kedro new --name=kedro_etl --tools=none --example=n 

Chuyển vào thư mục dự án. 

$ cd kedro-etl  

Tạo một thư mục trong thư mục pipelines tên data_processing.

$ mkdir -p src/kedro_etl/pipelines/data_processing  

Tạo file Python tên kedro_pipe.py và mở trong IDE bạn yêu thích, ví dụ Visual Studio Code.

$ code src/kedro_etl/pipelines/data_processing/kedro_pipe.py

Script Python nên chứa các hàm extract, transform và load, là các node trong pipeline. Ở đây lần lượt là các hàm create_sample_data(), clean_data(),load_and_process_data().

Sau đó, chúng ta nối các node bằng lớp Kedro Pipeline trong hàm create_pipeline(). Trong hàm pipeline, ta định nghĩa các node, mỗi node có inputs, outputs, và name của node. 

import pandas as pd
import numpy as np
from kedro.pipeline import Pipeline, node

def create_sample_data():
    data = {
        'id': range(1, 101),
        'name': [f'Person_{i}' for i in range(1, 101)],
        'age': np.random.randint(18, 80, 100),
        'salary': np.random.randint(20000, 100000, 100),
        'missing_values': [np.nan if i % 10 == 0 else i for i in range(100)]
    }
    return pd.DataFrame(data)

def clean_data(df: pd.DataFrame):
    # Remove rows with missing values
    df_cleaned = df.dropna()

    # Convert salary to thousands
    df_cleaned['salary'] = df_cleaned['salary'] / 1000

    # Capitalize names
    df_cleaned['name'] = df_cleaned['name'].str.upper()

    return df_cleaned

def load_and_process_data(df: pd.DataFrame):
    # Calculate average salary
    avg_salary = df['salary'].mean()

    # Add a new column for salary category
    df['salary_category'] = df['salary'].apply(
        lambda x: 'High' if x > avg_salary else 'Low')

    # Calculate age groups
    df['age_group'] = pd.cut(df['age'], bins=[0, 30, 50, 100], labels=[
                             'Young', 'Middle', 'Senior'])

    print(df)
    return df

def create_pipeline(**kwargs):
    return Pipeline(
        [
            node(
                func=create_sample_data,
                inputs=None,
                outputs="raw_data",
                name="create_sample_data_node",
            ),
            node(
                func=clean_data,
                inputs="raw_data",
                outputs="cleaned_data",
                name="clean_data_node",
            ),
            node(
                func=load_and_process_data,
                inputs="cleaned_data",
                outputs="processed_data",
                name="load_and_process_data_node",
            ),
        ]
    )

Nếu chạy pipeline mà không tạo data catalog, dữ liệu sẽ không được xuất. Vì vậy, chúng ta cần vào file conf/base/catalog.yml và chỉnh sửa, cung cấp cấu hình dataset.

raw_data:
  type: pandas.CSVDataset
  filepath: ./data/kedro/sample_data.csv

cleaned_data:
  type: pandas.CSVDataset
  filepath: ./data/kedro/cleaned_data.csv

processed_data:
  type: pandas.CSVDataset
  filepath: ./data/kedro/processed_data.csv

Chúng ta cũng phải đưa file Python vừa tạo vào registry của pipeline. Để làm vậy, vào file Python src/simple_etl/pipeline_registry.py và thêm đoạn mã sau. 

"""Project pipelines."""
from __future__ import annotations

from kedro.pipeline import Pipeline
from kedro_etl.pipelines.data_processing import kedro_pipe

def register_pipelines() -> Dict[str, Pipeline]:
    data_processing_pipeline = kedro_pipe.create_pipeline()

    return {
        "__default__": data_processing_pipeline,
        "data_processing": data_processing_pipeline,
    }

Chạy pipeline và xem log trực tiếp trong terminal bằng lệnh sau.

$ kedro run

Logs of Kedro pipeline run

Log lần chạy pipeline của Kedro.

Sau khi chạy pipeline, các file của bạn sẽ được lưu ở định dạng CSV tại vị trí đã định trong data catalog.

Output files of Kedro pipeline run

Các tệp đầu ra của lần chạy pipeline Kedro.

Nếu bạn gặp vấn đề khi chạy pipeline, hãy cân nhắc cài Kedro với đầy đủ phần mở rộng. 

$ pip install "kedro[all]"

Trực quan hóa Kedro

Chúng ta có thể trực quan hóa và chia sẻ pipeline bằng cách cài công cụ kedro-viz

$ pip install kedro-viz

Sau đó, chạy lệnh sau để trực quan hóa tất cả pipeline dữ liệu và các node dữ liệu. Công cụ cũng cung cấp tùy chọn truy vết thử nghiệm và khả năng chia sẻ trực quan hóa pipeline.

$ kedro viz run

Kedro Visualization

Trực quan hóa pipeline của Kedro.

5. Luigi

Luigi là framework bằng Python mã nguồn mở do Spotify phát triển, nổi trội trong quản lý các quy trình batch chạy dài và pipeline dữ liệu phức tạp. Công cụ mạnh về giải quyết phụ thuộc, quản lý workflow, trực quan hóa và phục hồi khi lỗi, giúp điều phối workflow dữ liệu hiệu quả. 

So với Airflow, Luigi có API tối giản, lập lịch theo lịch, và cộng đồng người dùng trung thành sẽ hỗ trợ bạn với mọi vấn đề liên quan đến pipeline điều phối dữ liệu. 

Nếu bạn mới học Python, có thể sẽ thấy khó để xây và chạy pipeline. Tuy nhiên, tài liệu và hướng dẫn có thể giúp bạn bắt đầu nhanh. Log cung cấp thông tin hạn chế, và dashboard chỉ là công cụ trực quan hóa DAG và phụ thuộc.

Abid Ali AwanAuthor

Bắt đầu với Luigi

Tạo pipeline dữ liệu với Luigi đòi hỏi hiểu biết về lập trình hướng đối tượng. Hãy bắt đầu bằng cách cài gói Python của Luigi. 

$ pip install luigi

Để phát triển pipeline ETL đơn giản trong Luigi, chúng ta sẽ tạo các tác vụ liên kết với nhau. Thay vì tạo các hàm Python làm tác vụ, chúng ta sẽ tạo một lớp Python cho mỗi bước trong pipeline, FetchData, ProcessDataGenerateReport. Mỗi lớp sẽ có ba hàm: requires(), output()run()

Hai hàm requires()output() sẽ kết nối các tác vụ, còn hàm run() sẽ thực thi mã xử lý. Cuối cùng, chúng ta sẽ xây pipeline bằng tác vụ cuối trong pipeline. 

import luigi
import pandas as pd
import numpy as np

class FetchData(luigi.Task):
    def output(self):
        return luigi.LocalTarget('data/fetch_data.csv')
    
    def run(self):
        # Simulate fetching data by creating a sample CSV file
        data = {
            'column1': [1, 2, np.nan, 4],
            'column2': ['A', 'B', 'C', np.nan]
        }
        df = pd.DataFrame(data)
        df.to_csv(self.output().path, index=False)

class ProcessData(luigi.Task):
    def requires(self):
        return FetchData()
    
    def output(self):
        return luigi.LocalTarget('data/process_data.csv')
    
    def run(self):
        df = pd.read_csv(self.input().path)
        # Fill missing values
        df['column1'].fillna(df['column1'].mean(), inplace=True)
        df['column2'].fillna('B', inplace=True)
        df.to_csv(self.output().path, index=False)

class GenerateReport(luigi.Task):
    def requires(self):
        return ProcessData()
    
    def output(self):
        return luigi.LocalTarget('data/generate_report.txt')
    
    def run(self):
        df = pd.read_csv(self.input().path)
        # Simple data analysis: calculate mean of column1 and value counts of column2
        mean_column1 = df['column1'].mean()
        value_counts_column2 = df['column2'].value_counts()
        
        with self.output().open('w') as out_file:
            out_file.write(f'Mean of column1: {mean_column1}\n')
            out_file.write('Value counts of column2:\n')
            out_file.write(value_counts_column2.to_string())

if __name__ == '__main__':
    luigi.build([GenerateReport()], local_scheduler=True)

Chạy đoạn mã trên trong Jupyter Notebook hoặc tạo file Python và chạy bằng terminal. 

Luigi Execution Summary

Tương tự như Luigi, bạn cũng có thể học cách xây dựng pipeline ETL với Apache Airflow. Hướng dẫn này bao quát những kiến thức cơ bản về trích xuất, biến đổi và nạp dữ liệu với Apache Airflow.

Luigi central planner

Chúng ta cần khởi tạo Luigi central planner để lập lịch chạy pipeline hoặc kích hoạt bằng sự kiện.

Khởi động scheduler bằng cách gõ lệnh sau trong terminal.

$ luigid  
2024-06-22 13:35:18,636 luigi[25056] INFO: logging configured by default settings
2024-06-22 13:35:18,636 luigi.scheduler[25056] INFO: No prior state file exists at /var/lib/luigi-server/state.pickle. Starting with empty state
2024-06-22 13:35:18,640 luigi.server[25056] INFO: Scheduler starting up

Để chạy pipeline, mở một terminal mới và gõ lệnh sau. Lệnh Luigi cần tên file Python và tác vụ cuối cùng mà ta muốn thực thi. Trong trường hợp này, tên file là luigi_pipe.py, và tác vụ Luigi cuối là GenerateReport.

$ python -m luigi --module luigi_pipe GenerateReport

Nếu bạn muốn trực quan hóa lần chạy pipeline và trạng thái tác vụ, chỉ cần truy cập http://localhost:8082 trong trình duyệt.

Luigi Central Planner webUI

Giao diện web Luigi Central Planner.

Như vậy là kết thúc phần hướng dẫn 5 lựa chọn thay thế tốt nhất cho Airflow! Nếu bạn muốn tìm hiểu sâu hơn bất kỳ ví dụ nào trong bài, dưới đây là một số tài nguyên nên tham khảo:

  • Đối với mã nguồn và dữ liệu của Prefect, Dagster và Luigi, vui lòng tham khảo không gian làm việc DataLab.
  • Đối với mã nguồn và dữ liệu của Mage AI và Kedro, vui lòng tham khảo kho GitHub.

Kết luận

Trong hướng dẫn này, chúng ta đã thảo luận các lựa chọn thay thế Airflow mã nguồn mở, miễn phí hàng đầu. Chúng ta cũng đã tìm hiểu từng công cụ điều phối dữ liệu, xây dựng và thực thi một pipeline ETL đơn giản. Những ví dụ mã này sẽ giúp bạn quyết định công cụ nào phù hợp nhất với trường hợp sử dụng của mình.

Nếu bạn là người mới bắt đầu, tôi khuyên nên bắt đầu với Prefect hoặc Mage AI vì chúng thân thiện với người dùng và cài đặt đơn giản. Tuy nhiên, nếu bạn tìm kiếm các công cụ nâng cao hơn và tuân thủ các thực hành kỹ nghệ phần mềm, hãy khám phá Dagster, Kedro và Luigi.

Sau khi đọc bài này, bước tiếp theo tự nhiên trong hành trình kỹ thuật dữ liệu của bạn là lấy chứng chỉ như Data Engineer in Python của DataCamp để học thêm các công cụ khác và xây dựng pipeline dữ liệu đầu-cuối có thể triển khai vào sản xuất.


Abid Ali Awan's photo
Author
Abid Ali Awan
LinkedIn
Twitter

Là một nhà khoa học dữ liệu được chứng nhận, tôi đam mê tận dụng công nghệ tiên tiến để tạo ra các ứng dụng học máy đổi mới. Với nền tảng vững chắc về nhận dạng giọng nói, phân tích và báo cáo dữ liệu, MLOps, AI hội thoại và NLP, tôi đã rèn giũa kỹ năng phát triển các hệ thống thông minh có thể tạo ra tác động thực sự. Bên cạnh chuyên môn kỹ thuật, tôi cũng là một người truyền đạt tốt, có khả năng chắt lọc các khái niệm phức tạp thành ngôn ngữ rõ ràng, súc tích. Nhờ đó, tôi trở thành một blogger được nhiều người quan tâm trong lĩnh vực khoa học dữ liệu, chia sẻ góc nhìn và kinh nghiệm với cộng đồng các chuyên gia dữ liệu ngày càng lớn. Hiện tại, tôi tập trung vào sáng tạo và biên tập nội dung, làm việc với các mô hình ngôn ngữ lớn để phát triển nội dung mạnh mẽ và hấp dẫn, giúp doanh nghiệp và cá nhân tận dụng tối đa dữ liệu của mình.

Chủ đề
Kỹ thuật Dữ liệu
Khoa học Dữ liệu

Tìm hiểu thêm về kỹ thuật dữ liệu với các khóa học này!

Courses

Introduction to Data Engineering

4 giờ
129.5K
Tìm hiểu về thế giới kỹ thuật dữ liệu trong khóa học ngắn này, bao gồm các công cụ và chủ đề như ETL và điện toán đám mây.
Xem chi tiếtRight Arrow
Bắt Đầu Khóa Học
Xem thêmRight Arrow
Có liên quan

blogs

Claude Opus 4.6: Tính năng, Điểm chuẩn, Bài kiểm tra thực hành và hơn thế nữa

Mô hình mới nhất của Anthropic dẫn đầu ở mã hóa tác tử và lập luận phức tạp. Thêm vào đó, nó có cửa sổ ngữ cảnh 1M.
Matt Crabtree's photo

Matt Crabtree

10 phút

Xem ThêmXem Thêm