1. Giới thiệu
Giả sử bạn là một Nhà khoa học dữ liệu tại Cymbal Financial, một đơn vị xử lý thanh toán có khối lượng giao dịch lớn. Đã xảy ra một loạt trường hợp chậm thanh toán và nhóm tuân thủ nghi ngờ có hành vi gian lận có phối hợp. Bạn cần tạo một quy trình để tiếp nhận nhật ký giao dịch thô của trung tâm thanh toán, dọn dẹp dữ liệu, huấn luyện một mô hình học máy, chạy suy luận hàng loạt và chuyển các giao dịch có rủi ro cao vào hàng đợi xem xét Cloud Spanner để kiểm tra thủ công.
Thông thường, việc này đòi hỏi bạn phải mất nhiều ngày để viết mã thiết lập lặp đi lặp lại (sổ tay Spark, cấu hình dbt, tập lệnh huấn luyện, DAG Airflow) và liên tục chuyển đổi ngữ cảnh giữa các giao diện và trình chỉnh sửa của bảng điều khiển.
Trong lớp học lập trình này, bạn sẽ lập trình theo cặp với một tác nhân bằng Bộ công cụ tác nhân dữ liệu (DAK) của Google Cloud trong IDE Antigravity. Bằng cách sử dụng ngôn ngữ tự nhiên đàm thoại, tác nhân này sẽ giúp bạn tạo sổ tay Spark, biên dịch một dự án dbt, tạo một vòng lặp suy luận và điều phối quy trình công việc bằng Dịch vụ được quản lý cho Apache Airflow.
Bạn sẽ thực hiện
- Nhập nhật ký của trung tâm thanh toán từ Cloud Storage bằng Dịch vụ được quản lý cho Apache Spark (Spark Serverless) vào một bảng BigQuery.
- Loại bỏ dữ liệu trùng lặp và chuẩn hoá giao dịch bằng cách sử dụng dbt để thiết lập các lớp dữ liệu sạch (Thô, Dàn dựng, Làm phong phú).
- Huấn luyện mô hình phân loại Rừng ngẫu nhiên phân tán (
RandomForestClassifier) trên Spark Serverless. - Chạy suy luận hàng loạt trên các giao dịch mới và ghi cảnh báo rủi ro cao trực tiếp vào Cloud Spanner.
- Điều phối, định cấu hình trực quan và triển khai toàn bộ quy trình bằng Dịch vụ được quản lý cho Apache Airflow và tính năng giám sát DAG tương tác trong IDE.
Bạn cần có
- Một trình duyệt web như Chrome
- Một dự án trên Google Cloud đã bật tính năng thanh toán (bạn nên sử dụng một dự án mới, chuyên dụng cho các lớp học thực hành).
- Hiểu biết cơ bản về SQL, Python và PySpark.
- Antigravity IDE khi có gói thuê bao Google AI Pro (nên dùng)
Các tài nguyên được tạo trong lớp học lập trình này sẽ có chi phí dưới 5 USD. Hãy nhớ làm theo hướng dẫn Dọn dẹp ở cuối phòng thí nghiệm để xoá các tài nguyên đã được cấp phép.
2. Thiết lập môi trường
Để bắt đầu phòng thí nghiệm, bạn sẽ chạy một tập lệnh khởi động. Tập lệnh này tự động bật các API GCP bắt buộc, tạo một bộ chứa Cloud Storage để nhập dữ liệu, tạo các tập dữ liệu giao dịch và thư mục mô phỏng, tải các thư mục tham chiếu vào BigQuery và bắt đầu cung cấp Cloud Spanner và Dịch vụ được quản lý cho Apache Airflow (trước đây gọi là Cloud Composer) ở chế độ nền.
Chọn hoặc tạo một dự án
Chọn một dự án hiện có hoặc tạo một dự án mới trong Google Cloud Console.
Xác minh thông tin thanh toán
Đảm bảo rằng bạn đã bật tính năng thanh toán cho dự án trên đám mây của mình trên Google Cloud. Bạn có thể tìm hiểu thêm về cách thực hiện việc này bằng cách làm theo hướng dẫn này.
Chạy tập lệnh thiết lập
Bạn sẽ sử dụng Google Cloud Shell (hoặc shell cục bộ được định cấu hình bằng Google Cloud CLI) để khởi chạy quá trình thiết lập môi trường.
- Mở Google Cloud Console.
- Nhấp vào Kích hoạt Cloud Shell trong thanh công cụ ở trên cùng bên phải.

- Trong thiết bị đầu cuối Cloud Shell, hãy định cấu hình dự án đang hoạt động của bạn:
gcloud config set project <<YOUR_PROJECT_ID>>
export PROJECT_ID=$(gcloud config get-value project)
- Sao chép kho lưu trữ lớp học lập trình và chuyển đến thư mục tập lệnh:
cd ~/
git clone --filter=blob:none --no-checkout https://github.com/GoogleCloudPlatform/devrel-demos.git
cd ~/devrel-demos
git sparse-checkout init --cone
git sparse-checkout set codelabs/agentic-data-labs/data-science
git checkout main
cd codelabs/agentic-data-labs/data-science/scripts
- Chạy tập lệnh thiết lập khởi động để triển khai tất cả tài nguyên vào
us-central1:
chmod +x setup.sh setup_spanner.sh setup_composer.sh
export REGION=us-central1
./setup.sh
- Khi tập lệnh hoàn tất, bạn sẽ thấy một bản tóm tắt cho biết tập dữ liệu BigQuery và bộ chứa Cloud Storage của bạn đã sẵn sàng. Trong nền, Cloud Spanner (mất khoảng 2 phút) và Managed Airflow (mất khoảng 20 phút) sẽ tiếp tục cung cấp. Bạn có thể theo dõi tiến trình của họ bất cứ lúc nào bằng cách chạy:
tail -f /tmp/spanner_setup.log
tail -f /tmp/composer_setup.log
Mở Antigravity IDE
- Tải và cài đặt Antigravity IDE từ trang tải xuống Google Antigravity.
- Khởi chạy Antigravity IDE.
- Tạo một thư mục mới, trống trên máy cục bộ (ví dụ: đặt tên là
agentic-data-labs) rồi mở thư mục đó trong IDE bằng cách chọn Open Folder (Mở thư mục). Đây sẽ là không gian làm việc cục bộ của bạn cho lớp học lập trình.

Cài đặt tiện ích Data Agent Kit
Tiện ích Google Cloud Data Agent Kit cung cấp khả năng tích hợp sâu với các dịch vụ dữ liệu của Google Cloud ngay trong trình chỉnh sửa, cho phép bạn tương tác với BigQuery, Cloud SQL, Cloud Storage và nhiều dịch vụ khác mà không cần chuyển đổi ngữ cảnh.
- Trong Antigravity IDE, hãy nhấp vào biểu tượng Tiện ích trong Thanh hoạt động ở ngoài cùng bên trái màn hình (biểu tượng này trông giống như 4 hình vuông).
- Trong thanh tìm kiếm ở đầu ngăn Tiện ích, hãy nhập
Google Cloud Data Agent Kit. - Tìm tiện ích có tên Google Cloud Data Agent Kit do
googlecloudtoolsxuất bản - Nhấp vào nút Install (Cài đặt).
- Một lời nhắc có thể xuất hiện và hỏi "Bạn có tin tưởng nhà xuất bản "googlecloudtools" và các tiện ích của họ không?". Nhấp vào Tin tưởng nhà xuất bản và cài đặt để tiếp tục.

Sau khi cài đặt, bạn sẽ thấy biểu tượng Google Cloud Data Agent Kit mới xuất hiện trong Thanh hoạt động ở phía bên trái của Antigravity IDE.
- Một trang giới thiệu có tiêu đề "Chào mừng bạn đến với Google Cloud Data Agent Kit" sẽ tự động mở ra. Nếu bạn chưa đăng nhập vào tài khoản Cloud, hãy làm theo mọi lời nhắc để cho phép truy cập.
- Trong phần Configuration Summary (Tóm tắt cấu hình), hãy tìm trường dự án. Nhấp vào trình đơn thả xuống rồi chọn dự án trên đám mây của bạn trên Google Cloud. Đặt khu vực của bạn là
us-central1. Sau đó, chọn Configure MCP Servers (Định cấu hình máy chủ MCP).

- Chọn Định cấu hình máy chủ MCP. Trong ngăn MCP Configuration (Cấu hình MCP), hãy đảm bảo bạn bật các máy chủ MCP từ xa sau đây:
- BigQuery
- Spanner
- Sổ tay
Sau đó, hãy nhấp vào Bắt đầu.

Khám phá các lựa chọn về cấu hình
Sau khi hoàn tất quá trình thiết lập, bạn sẽ được chuyển đến trang "Bắt đầu sử dụng Bộ công cụ Data Agent của Google Cloud".
- Trong mục "Thiết lập và cấu hình", hãy nhấp vào Bắt đầu.
- Thao tác này sẽ mở bảng điều khiển Cấu hình bộ công cụ Data Agent. Khám phá các thẻ:
- Dự án và khu vực: Xác minh mã dự án bạn đã chọn và xác nhận rằng tập lệnh thiết lập đã bật tất cả các API bắt buộc (Compute Engine, Cloud Storage, BigQuery, Spanner, v.v.).
- BigQuery: Định cấu hình vị trí mặc định cho các truy vấn BigQuery. Sử dụng khu vực
us-central1. - Định cấu hình máy chủ MCP: Xem các máy chủ MCP (BigQuery, Notebooks, Spanner, v.v.) đã bật để cho phép các tác nhân AI tương tác an toàn với dữ liệu của bạn.
- Kỹ năng: Khám phá các kỹ năng được tạo sẵn giúp cung cấp cho các tác nhân những khả năng chuyên biệt để thực hiện các tác vụ dữ liệu phức tạp.

Tóm tắt phần: Bạn đã chạy tập lệnh khởi động để tạo thành phần GCS và BigQuery trong khi Spanner và Airflow được tạo ở chế độ nền. Sau đó, bạn mở dự án trong Antigravity IDE và kích hoạt tiện ích Google Cloud Data Agent Kit. Giờ đây, bạn đã sẵn sàng viết sổ tay đầu tiên.
3. Nhập nhật ký thô bằng Spark Serverless
Trong phần này, bạn sẽ nhập nhật ký giao dịch JSON thô vào hồ dữ liệu. Dịch vụ được quản lý cho Apache Spark (Spark Serverless) kết nối trực tiếp với bộ nhớ gốc của BigQuery. Bạn sẽ sử dụng trình kết nối BigQuery tiêu chuẩn để quản lý dữ liệu dạng bảng và cho phép truy vấn cũng như phân tích trực tiếp.
Khám phá thời gian chạy Spark Serverless được định cấu hình sẵn
Trước khi thực thi mã Spark, hãy kiểm tra mẫu Thời gian chạy không máy chủ do tập lệnh thiết lập định cấu hình trước. Mẫu này xác định phần phụ trợ của môi trường thực thi mục tiêu và gói các phần phụ thuộc cần thiết của trình kết nối.
- Trong thanh hoạt động của IDE, hãy mở bảng điều khiển Google Cloud Data Agent Kit.
- Mở rộng trình đơn thả xuống Apache Spark, sau đó mở rộng Serverless.
- Nhấp chuột phải vào
fraud-pipeline-runtimerồi chọn Profile (Hồ sơ) để mở chế độ xem cấu hình của hồ sơ đó trong trình chỉnh sửa. - Trong thẻ Profile (Hồ sơ), hãy di chuyển xuống rồi mở rộng Properties (Thuộc tính) để kiểm tra các phần phụ thuộc tuỳ chỉnh được đính kèm vào môi trường:
spark.jars: Chứags://spark-lib/spanner/spark-3.5-spanner-1.4.0.jar, sử dụng trình kết nối Spark Spanner để cho phép các tác vụ Spark ghi kết quả suy luận trực tiếp vào Cloud Spanner sau này trong phòng thí nghiệm. (Lưu ý: Dataproc Serverless bao gồm trình kết nối Spark BigQuery của Google Cloud theo mặc định, không yêu cầu cấu hình jar bổ sung để đọc và ghi các bảng BigQuery).

- Hãy lưu ý thẻ Phiên tương tác ở bên trái. Hiện tại, bảng này trống vì bạn chưa thực thi mã nào. Ngay khi bạn chạy sổ tay trong bước tiếp theo, một phiên tính toán trực tiếp không cần máy chủ sẽ được cung cấp động và xuất hiện tại đây!
Nhập dữ liệu bằng Bộ công cụ tác nhân dữ liệu
Thay vì định cấu hình một Phiên Spark theo cách thủ công hoặc viết các tập lệnh tải PySpark từ đầu, bạn sẽ lập trình theo cặp với một tác nhân bằng cách sử dụng Data Agent Kit.
- Mở ngăn Trò chuyện với nhân viên hỗ trợ bằng cách nhấp vào biểu tượng Bật/tắt nhân viên hỗ trợ trong thanh công cụ trên cùng bên phải.
- Dán câu lệnh sau vào cuộc trò chuyện (nhớ thay thế
${PROJECT_ID}bằng mã dự án thực tế của bạn trên Google Cloud):
Create a PySpark notebook (01_ingestion.ipynb) to ingest JSON transaction logs
from gs://${PROJECT_ID}-fin-clearing-raw/ into a BigQuery table
`${PROJECT_ID}.transactions_dataset_evals.raw_transactions`
using the Spark BigQuery connector (`format("bigquery")`) with overwrite mode.
- Nếu nhân viên hỗ trợ yêu cầu quyền thực thi các lệnh xác minh ở chế độ nền (ví dụ: "Cho phép chạy lệnh này?"), hãy xem xét lệnh được đề xuất rồi chọn Có, cho phép lần này (hoặc Có và luôn cho phép).
- Khi tác nhân hoàn tất việc tạo tệp, hãy nhấp vào nút Chấp nhận tất cả (hoặc biểu tượng dấu đánh dấu) màu xanh dương ở cuối ngăn trò chuyện để lưu
notebooks/01_ingestion.ipynbvào không gian làm việc của bạn.

Xem xét và thực thi sổ tay
- Mở
notebooks/01_ingestion.ipynbmới tạo trong IDE. - Xem lại mã PySpark cho logic ghi của trình kết nối BigQuery.
- Nhấp vào Run All (Chạy tất cả) trong thanh công cụ sổ tay của IDE.
- Nếu đây là lần đầu tiên bạn chạy một sổ tay Spark từ xa, IDE có thể nhắc bạn cài đặt các phần phụ thuộc cục bộ. Nếu được nhắc, hãy nhấp vào Install dependencies for Remote Spark Kernels (Cài đặt các phần phụ thuộc cho Nhân Spark từ xa) và xác nhận các hộp thoại cài đặt, sau đó nhấp lại vào Run All (Chạy tất cả).
- Trong trình đơn thả xuống Select Kernel (Chọn hạt nhân), hãy chọn Remote Spark Kernels (Hạt nhân Spark từ xa) -> fraud-pipeline-runtime on Serverless Spark (fraud-pipeline-runtime trên Serverless Spark). (Lưu ý: Nếu bạn không thấy mẫu thời gian chạy được định cấu hình trước trong danh sách, hãy nhấp vào biểu tượng làm mới ở trên cùng bên phải của trình chọn nhân để tải lại các nhân từ xa hiện có).
- Nhìn vào thanh trạng thái ở dưới cùng bên trái của trình chỉnh sửa. Bạn sẽ thấy
Connecting to kernel: fraud-pipeline-runtime on Serverless Spark.... Vì đây là lần đầu tiên ra mắt phần phụ trợ hạt nhân thời gian chạy Spark Serverless, nên sẽ mất vài phút để cung cấp và khởi động. - Sau khi kernel hoàn tất việc kết nối, sổ tay sẽ tự động bắt đầu thực thi tất cả các ô theo trình tự để xử lý nhật ký giao dịch thô thành tập dữ liệu BigQuery của bạn.
Xác minh
Sau khi quá trình thực thi hoàn tất, hãy kiểm tra danh mục Data Agent Kit để xác minh việc tạo bảng:

- Trong thanh hoạt động của IDE, hãy mở bảng điều khiển Google Cloud Data Agent Kit.
- Mở rộng mục CATALOG.
- Mở rộng mã dự án của bạn.
- Mở rộng BigQuery.
- Mở rộng tập dữ liệu
transactions_dataset_evals. - Nhấp vào bảng
raw_transactionsđể mở chế độ xem chi tiết của bảng trong trình chỉnh sửa chính. - Trong bảng điều hướng bên trái, hãy khám phá các thẻ Dữ liệu, Giản đồ và Chi tiết để kiểm tra các bản ghi và siêu dữ liệu đã được nhập.
Tóm tắt phần: Bạn đã sử dụng ngôn ngữ tự nhiên trong Agent Chat để tạo một khối lượng công việc hoàn chỉnh cho Spark Serverless. Sau đó, bạn đã thực thi để xử lý nhật ký JSON không có cấu trúc thành một bảng BigQuery (thô).
4. Loại bỏ dữ liệu trùng lặp và chuẩn hoá bằng dbt
Trước khi huấn luyện mô hình học máy, bạn sẽ thực thi chất lượng dữ liệu bằng cách xoá các nhật ký phát trực tuyến trùng lặp, tách biệt các bản ghi không hợp lệ (chẳng hạn như mã giao dịch trống) và kết hợp dữ liệu theo phương diện (người thanh toán và người nhận thanh toán). Quy trình này yêu cầu các hoạt động chuyển đổi SQL đáng tin cậy và có tính chất luỹ đẳng, khiến dbt (công cụ tạo dữ liệu) trở thành một lựa chọn phù hợp.
Tạo giàn giáo cho quy trình dbt
Sử dụng tác nhân để tạo một dự án dbt trên tập dữ liệu BigQuery:
- Quay lại ngăn Agent Chat (Trò chuyện với nhân viên hỗ trợ).
- Hãy cung cấp hướng dẫn sau để tạo dự án dbt:
Scaffold a us-central1 dbt project in dbt_project/ that maps raw_transactions
through to an enriched_transactions model in dataset transactions_dataset_evals.
Deduplicate by transaction_id in staging. Quarantine null IDs to an invalid_transactions model.
Join the valid staging records with dim_payers and dim_payees for the enriched_transactions model,
preserving the historical `is_fraud` label column, and finally add a transaction uniqueness test.
Create an implementation plan first.
- Tác nhân sẽ trình bày một cấu phần phần mềm Kế hoạch triển khai trong ngăn trình chỉnh sửa chính. Xem xét cấu trúc tệp và logic SQL được đề xuất.
- Nhấp vào Tiếp tục (rồi nhấp vào Chấp nhận tất cả) để cho phép tác nhân tạo tệp trong không gian làm việc của bạn.

- Sau khi quá trình tạo hoàn tất, tác nhân sẽ hiển thị một Hướng dẫn tóm tắt các thành phần mới. Chấp nhận tất cả các thay đổi nếu được nhắc.

Xây dựng và thử nghiệm
Mặc dù tác nhân tự động chạy dbt compile để đảm bảo SQL được tạo hợp lệ về mặt cú pháp, nhưng giờ đây, bạn sẽ hiện thực hoá các khung hiển thị và bảng này vào BigQuery, đồng thời chạy các kiểm thử chất lượng dữ liệu để xác minh cục bộ. (Lưu ý: Sau này trong phòng thí nghiệm, bạn sẽ tự động hoá bước dbt này trong DAG Airflow toàn diện).
- Trong thanh hoạt động ở ngoài cùng bên trái, hãy nhấp vào biểu tượng Trình khám phá (hoặc nhấn
Cmd/Ctrl+Shift+E). - Mở rộng
dbt_project->modelsđể kiểm tra các mô hình SQL được tạo. Nhấp vàoenriched_transactions.sqlđể mở và xem lại logic của tính năng biến đổi và phát hiện hành vi gian lận trong trình chỉnh sửa. - Trong File Explorer, hãy nhấp chuột phải vào thư mục
dbt_projectrồi chọn Open in Integrated Terminal (Mở trong Integrated Terminal). Thao tác này sẽ tự động mở một ngăn dòng lệnh được đặt trực tiếp vào thư mục làm việcdbt_projectbắt buộc. - Nếu bạn chưa cài đặt
dbt, hãy tạo một môi trường ảo bên ngoàidbt_project/(tại thư mục gốc của nhà hoặc không gian làm việc) rồi cài đặt trình kết nối BigQuery:
python3 -m venv ~/.venv/dbt
source ~/.venv/dbt/bin/activate
pip install dbt-bigquery
- Chạy các mô hình dbt và các kiểm thử chất lượng dữ liệu liên quan:
dbt build
- Xem đầu ra của thiết bị đầu cuối. dbt sẽ biên dịch SQL, hiện thực hoá các bảng dàn dựng và bảng được làm phong phú trong BigQuery, đồng thời thực thi các kiểm thử dữ liệu.

- Sau khi quá trình tạo hoàn tất, hãy đóng ngăn cửa sổ để giải phóng không gian màn hình cho các bước còn lại.
Tóm tắt phần: Bạn đã tạo một dự án dbt bằng tác nhân, chạy các kiểm thử chất lượng dữ liệu và chuyển đổi các bản ghi thô thành các bảng BigQuery dàn dựng và làm phong phú.
5. Huấn luyện mô hình phát hiện gian lận phân tán bằng Random Forest
Với các giao dịch được làm giàu và cụ thể hoá trong BigQuery, bạn sẽ xây dựng một mô hình học máy để phân loại các sự kiện gian lận. Rừng ngẫu nhiên là một phương pháp học tập kết hợp phù hợp với dữ liệu phân loại dạng bảng. Việc chạy RandomForestClassifier trên Spark Serverless sẽ phân phối quá trình huấn luyện mô hình trên các nút worker mà không yêu cầu bạn quản lý cơ sở hạ tầng.
Trong bước này, bạn sẽ dùng tác nhân để tạo quy trình huấn luyện Spark ML.
Tạo sổ tay huấn luyện ML
- Mở ngăn Chat với nhân viên hỗ trợ.
- Đưa ra câu lệnh sau để thiết kế trình tự huấn luyện mô hình (hãy nhớ thay thế
${PROJECT_ID}bằng mã dự án đang hoạt động của bạn):
Create a PySpark notebook (02_training.ipynb) to train a distributed Random Forest
(RandomForestClassifier) model on the BigQuery table
`transactions_dataset_evals`.`enriched_transactions`, predicting the `is_fraud` label.
Train only on historically labeled records where `is_fraud` is not null.
One-hot encode categorical strings, scale amounts, cache the dataset in memory,
evaluate AUC, and save the evaluated model to gs://${PROJECT_ID}-models/fraud_model.
- Xem lại kế hoạch hoặc mã được tạo của tác nhân rồi nhấp vào Tiến hành / Chấp nhận tất cả để lưu
notebooks/02_training.ipynbvào không gian làm việc của bạn.

Xem xét và thực thi sổ tay
- Mở
notebooks/02_training.ipynbtrong trình chỉnh sửa. - Xem xét các giai đoạn của quy trình ML PySpark để mã hoá đối tượng, tập hợp vectơ và logic phân loại Rừng ngẫu nhiên.
- Nhấp vào Run All (Chạy tất cả) trong thanh công cụ sổ tay của IDE.
- Khi trình chọn thả xuống Select Kernel (Chọn nhân) mở ra, hãy chọn fraud-pipeline-runtime on Serverless Spark (thời gian chạy quy trình phát hiện gian lận trên Serverless Spark).

Xác minh
Sau khi quá trình thực thi hoàn tất, hãy xác nhận rằng mô hình đã được huấn luyện và xuất đúng cách:
- Xem xét các đầu ra của ô đánh giá ở gần cuối sổ tay để xác minh Điểm diện tích dưới đường cong ROC (AUC) được báo cáo.
- Để đảm bảo các cấu phần phần mềm của mô hình được lưu thành công vào GCS, hãy mở rộng ngăn trình khám phá STORAGE (LƯU TRỮ) trong thanh bên Data Agent Kit.
- Xác định vị trí của nhóm kết thúc bằng
-models(được liên kết với Mã dự án đang hoạt động của bạn), mở rộng nhóm đó và đi sâu vào để xác minh rằng thư mụcfraud_modelvà các giai đoạn của quy trình tồn tại.

Tóm tắt phần: Bạn đã dùng tác nhân để tạo quy trình huấn luyện PySpark ML, huấn luyện mô hình Rừng ngẫu nhiên trên bảng BigQuery được làm phong phú và xuất mô hình sang Cloud Storage.
6. Suy luận hàng loạt và ghi Cloud Spanner
Với một mô hình dự đoán đã được huấn luyện và lưu trữ trong Cloud Storage, bạn sẽ chạy suy luận theo lô trên các giao dịch mới diễn ra thông qua BigQuery. Các giao dịch có rủi ro cao cần được chuyển đến một hệ thống vận hành để nhóm tuân thủ có thể xem xét. Cloud Spanner cung cấp một cơ sở dữ liệu giao dịch có khả năng mở rộng cho hàng đợi đánh giá này.
Tạo sổ tay suy luận theo lô
Sử dụng tác nhân để tạo một sổ tay suy luận kết nối BigQuery, Cloud Storage và Cloud Spanner:
- Mở ngăn Chat với nhân viên hỗ trợ.
- Cung cấp câu lệnh sau:
Create an inference notebook (03_inference.ipynb) that loads the RandomForestClassifier
model to score unlabeled records (where `is_fraud` is null) from the BigQuery table
`transactions_dataset_evals`.`enriched_transactions`.
Filter for high-risk transactions with a 50%+ fraud probability score (probability >= 0.50)
and write them to the Cloud Spanner table SparkEvalFraudReviewQueue in the cymbal-fraud instance
under fraud-db.
- Chấp nhận sổ tay được tạo để lưu
notebooks/03_inference.ipynbvào không gian làm việc của bạn.

Xem xét và thực thi sổ tay
- Mở
notebooks/03_inference.ipynbvừa tạo trong trình chỉnh sửa. - Xem lại trình tự suy luận PySpark:
- Phần phụ thuộc: Mẫu Thời gian chạy phi máy chủ cung cấp các phần phụ thuộc JAR
cloud-spannerbắt buộc để thực thi Spark. - Định dạng dữ liệu: Tập lệnh sẽ loại bỏ các cột vectơ Spark ML phức tạp (chẳng hạn như các đặc điểm và xác suất thô) trước khi ghi để khớp với giản đồ bảng Spanner.
- Trình kết nối Spanner: Trình kết nối này ghi các hàng được gắn cờ bằng cách sử dụng
.format("cloud-spanner")để thêm trực tiếp vào hàng đợi đánh giá.
- Phần phụ thuộc: Mẫu Thời gian chạy phi máy chủ cung cấp các phần phụ thuộc JAR
- Nhấp vào Run All (Chạy tất cả) trong thanh công cụ sổ tay của IDE.
- Khi được nhắc chọn một hạt nhân, hãy chọn fraud-pipeline-runtime trên Serverless Spark.
Xác minh
Sau khi sổ tay suy luận hoàn tất quá trình xử lý, bạn có thể truy vấn cơ sở dữ liệu Spanner hoạt động ngay trong IDE:
- Trong thanh hoạt động của IDE, hãy mở bảng điều khiển Google Cloud Data Agent Kit.
- Mở rộng mục CATALOG.
- Mở rộng mã dự án của bạn, sau đó mở rộng Spanner.
- Chuyển đến
cymbal-fraud->fraud-db->Tables->SparkEvalFraudReviewQueue. - Nhấp chuột phải vào bảng rồi chọn Bảng truy vấn, sau đó thực thi truy vấn:
SELECT *
FROM `SparkEvalFraudReviewQueue`
LIMIT 100;
- Trong ngăn Query Results (Kết quả truy vấn) bên dưới, bạn sẽ thấy các hàng mới được chèn đại diện cho những giao dịch có rủi ro cao được gắn cờ để xem xét theo cách thủ công.

Tóm tắt phần: Bạn đã dùng tác nhân để tạo một sổ tay suy luận theo lô, tính điểm cho các bản ghi BigQuery chưa được gắn nhãn bằng mô hình đã huấn luyện và ghi các giao dịch có rủi ro cao trực tiếp vào Cloud Spanner.
7. Tạo khung và điều phối bằng Managed Airflow
Hiện tại, quy trình của bạn bao gồm các bước riêng biệt: sổ tay tiếp nhận, dự án chuyển đổi dbt và sổ tay suy luận theo lô. Để chuẩn bị cho quá trình sản xuất, bạn sẽ ghép các ô này lại với nhau thành một biểu đồ phụ thuộc theo lịch.
Dịch vụ được quản lý cho Apache Airflow (trước đây là Cloud Composer) cung cấp một công cụ điều phối được quản lý cho quy trình công việc này. Data Agent Kit có tính năng Orchestration Pipelines (Các quy trình phối hợp) giúp chuyển đổi trực tiếp các định nghĩa quy trình YAML khai báo thành DAG Airflow.
Xác định quy trình
Sử dụng tác nhân để tạo cấu hình quy trình điều phối:
- Trong Agent Chat (Trò chuyện với tác nhân), hãy đưa ra câu lệnh sau (nhớ thay thế
${PROJECT_ID}):
Initialize and define an orchestration pipeline (fraud_analysis_pipeline) triggering
the ingestion notebook, dbt project, and inference notebook in sequential order.
For Dataproc Serverless engine configs in us-central1, use resourceProfile.inline
(defining properties with spark.jars: "gs://spark-lib/spanner/spark-3.5-spanner-1.4.0.jar"
for inference) rather than resourceProfile.path or overrides.
Set the schedule interval to run daily at midnight, and use
gs://${PROJECT_ID}-airflow-artifacts for artifact storage.
Xem xét cấu hình DAG
Trình điều phối Bộ công cụ tác nhân dữ liệu sử dụng các cấu hình YAML khai báo để xác định và triển khai các quy trình đến Apache Airflow, cho phép kiểm soát phiên bản và triển khai các định nghĩa thông qua CI/CD.
Trong ngăn IDE Explorer, hãy xem xét 2 tệp quy trình mà tác nhân đã tạo ở gốc của không gian làm việc:
deployment.yaml: Mở tệp này. Đây là sổ đăng ký môi trường của bạn. Nó liên kết quy trìnhdevlogic của bạn với môi trườngcymbal-airflow, đặt khu vực thực thi (us-central1) và xác định nhómartifact_storagenơi các DAG và phần phụ thuộc đã biên dịch được dàn dựng.fraud_analysis_pipeline.yaml: Mở tệp này. Thao tác này xác định biểu đồ thực thi. Thao tác này chỉ định lịch kích hoạt (interval: '0 0 * * *') và sắp xếp theo trình tự 3 bước trong khốiactions:- Một thao tác
notebooktiếp nhận cho01_ingestion.ipynbđang chạy trên Dataproc Serverless. - Một thao tác chuyển đổi
pipelinenhắm đến thư mụcdbt_project, với một phần phụ thuộcdependsOntrỏ đến bước tiếp nhận. - Một thao tác suy luận
notebookcho03_inference.ipynbvới một phần phụ thuộcdependsOntrỏ đến bước dbt, liên kết thuộc tính JAR của Spanner.
- Một thao tác
- Tác nhân này cũng sẽ tóm tắt những cấu phần phần mềm đã tạo này vào thẻ Hướng dẫn trong ngăn trình chỉnh sửa, nêu rõ các cấu hình và hoạt động xác thực đã thực hiện.
Cấu hình DAG tương tác
Data Agent Kit hiển thị cấu hình quy trình dưới dạng biểu đồ trực quan tương tác để kiểm tra và chỉnh sửa các thuộc tính DAG của Airflow.
- Trong thanh hoạt động của IDE, hãy mở bảng điều khiển Google Cloud Data Agent Kit.
- Trong phần
DATA ENGINEERING, hãy mở rộngOrchestration Pipelines. - Nhấp vào
fraud_analysis_pipeline.yamlđể mở canvas DAG trực quan trong trình chỉnh sửa chính.

- Nhấp vào nút
Schedule triggerở trên cùng. Một bảng điều khiển cấu hình sẽ mở ra ở bên phải, hiển thị chuỗi Cron đã phân tích cú pháp (0 0 * * *) và cho phép bạn điều chỉnh các tham số như backfill và catchup. - Nhấp vào nút tác vụ sổ tay (chẳng hạn như bước truyền dữ liệu hoặc suy luận). Flyout sẽ cập nhật để hiển thị các ánh xạ thực thi Dataproc Serverless cụ thể và các thuộc tính của trình kết nối.
- Lưu ý siêu liên kết tên tệp sổ tay (chẳng hạn như
01_ingestion.ipynb) bên trong khối nút. Khi bạn nhấp vào biểu tượng này, sổ tay sẽ mở ra ngay trong trình chỉnh sửa. - Trong thanh bên trái bên dưới Orchestration Pipelines (Các quy trình phối hợp), hãy nhấp vào
Deployment configuration. Chế độ xem này cho thấy cụm môi trườngdevmục tiêu và các cấu phần phần mềm của bộ chứa GCS đầu ra.
Tóm tắt phần: Bạn đã tạo một cấu hình quy trình điều phối bằng tác nhân, xác định các phần phụ thuộc giữa các tác vụ tiếp nhận, dbt và suy luận trong một canvas trực quan tương tác.
8. Triển khai, thực thi và giám sát
Khi DAG được xác định cục bộ, bạn sẽ kết nối với môi trường Airflow được quản lý đã được cung cấp trong quá trình thiết lập và triển khai quy trình.
Định cấu hình Dịch vụ được quản lý cho Apache Airflow
Trước khi triển khai, hãy định cấu hình mối kết nối Trình lập lịch trong phần cài đặt Data Agent Kit để tiện ích nhắm đến môi trường Airflow được quản lý của bạn:
- Trong thanh hoạt động của IDE, hãy mở bảng điều khiển Google Cloud Data Agent Kit.
- Trong phần
SETTINGS, hãy nhấp vào Cài đặt. - Chọn Trình lập lịch biểu trong trình đơn bên trái.
- Định cấu hình chế độ cài đặt:
- Mã dự án: Chọn mã dự án đang hoạt động.
- Khu vực: Chọn
us-central1. - Môi trường: Chọn
cymbal-airflow.
- Nhấp vào Lưu.

Triển khai DAG
Giờ đây, bạn sẽ triển khai quy trình được định cấu hình trực tiếp vào môi trường Airflow được quản lý từ canvas trực quan:
- Trong thanh bên Google Cloud Data Agent Kit, hãy mở rộng
DATA ENGINEERING>Orchestration Pipelinesrồi nhấp vàofraud_analysis_pipeline.yamlđể mở canvas DAG trực quan. - Ở góc trên cùng bên phải của thanh công cụ canvas, hãy nhấp vào nút Chạy quy trình màu xanh dương.
- Trong trình chọn môi trường thả xuống, hãy chọn
dev. - Theo dõi thông báo tiến trình ở khu vực trạng thái dưới cùng (
Running pipeline: Building pipeline locally...). Tiện ích này sẽ tự động biên dịch DAG, đóng gói sổ tay và các thành phần dbt, rồi tải chúng lên bộ chứa GCS của môi trường Airflow được quản lý (quá trình này mất khoảng 3 đến 4 phút để hoàn tất).

Theo dõi lượt chạy
Sau khi quá trình biên dịch cục bộ hoàn tất và thông báo bật lên xác nhận Triggered a new run for pipeline... successfully, hãy theo dõi quá trình thực thi trực tiếp:
- Trong thanh bên Google Cloud Data Agent Kit, hãy mở rộng
DATA ENGINEERING>Orchestration Pipelines. - Nhấp vào Quản lý quy trình.
- Trong bảng Quản lý quy trình, hãy nhấp vào biểu tượng
fraud_analysis_pipelineđể mở nhật ký thực thi của quy trình.

- Trong chế độ xem Execution History (Nhật ký thực thi), hãy chọn lần chạy đang hoạt động trong lịch.
- Khi quá trình thực thi diễn ra trên từng tác vụ trong quy trình (truy nạp, chuyển đổi dbt và suy luận), các chỉ báo trạng thái sẽ cập nhật và thời lượng tác vụ sẽ được điền sẵn. Nhấp vào một tác vụ bất kỳ để kiểm tra đầu ra thực thi trực tiếp và nhật ký DAG của Airflow.

Tóm tắt phần: Bạn đã định cấu hình kết nối Airflow Scheduler, triển khai quy trình phân tích toàn diện cho Managed Airflow và giám sát quá trình thực thi trực tiếp, xác minh hệ thống từ nhật ký thô đến các dự đoán cuối cùng của Cloud Spanner.
9. Dọn dẹp
Để tránh bị tính phí liên tục cho dự án trên đám mây của bạn đối với các tài nguyên được dùng trong lớp học lập trình này, hãy huỷ môi trường bằng cách sử dụng tập lệnh tự động.
- Trong bảng điều khiển Terminal (hoặc trong Cloud Shell), hãy chuyển đến thư mục tập lệnh rồi thực thi:
cd ~/devrel-demos/codelabs/agentic-data-labs/data-science/scripts
chmod +x teardown.sh
./teardown.sh
- Tập lệnh sẽ liệt kê tất cả các tài nguyên mà tập lệnh dự định xoá và nhắc bạn xác nhận:
- Môi trường Airflow được quản lý (
cymbal-airflow) - Cloud Spanner Instance (
cymbal-fraud) - Tập dữ liệu BigQuery (
transactions_dataset_evals) - Bộ chứa Cloud Storage (
gs://${PROJECT_ID}-fin-clearing-rawvàgs://${PROJECT_ID}-models) - Tài khoản dịch vụ của nhân viên (
composer-worker-sa)
- Môi trường Airflow được quản lý (
- Nhập
yđể xác nhận. Tập lệnh huỷ sẽ xoá tất cả các dịch vụ GCP được cung cấp và dọn dẹp các tệp cục bộ.
10. Xin chúc mừng!
Bạn đã xây dựng một quy trình phát hiện gian lận toàn diện trải rộng trên Cloud Storage, BigQuery, Managed Service for Apache Spark (Spark Serverless), dbt, Cloud Spanner và Managed Service for Apache Airflow, lập trình theo cặp với Google Cloud Data Agent Kit trong Antigravity IDE.
Thành tích bạn đạt được
- 📥 Đã chuyển nhật ký giao dịch thô vào một bảng BigQuery bằng Managed Service for Apache Spark và Data Agent Kit.
- 🧹 Loại bỏ dữ liệu trùng lặp và chuẩn hoá dữ liệu bằng cách tạo một dự án dbt có các kiểm thử chất lượng dữ liệu.
- 🤖 Huấn luyện mô hình Rừng ngẫu nhiên phân tán bằng cách sử dụng
RandomForestClassifiervà xuất mô hình đã huấn luyện sang Cloud Storage. - ⚡ Thực hiện suy luận hàng loạt trên các giao dịch đến và chuyển các bản ghi có rủi ro cao vào Cloud Spanner để xem xét kiểm tra.
- 🔄 Điều phối, triển khai và giám sát quy trình công việc dưới dạng DAG Airflow theo lịch bằng cách sử dụng Dịch vụ được quản lý cho Apache Airflow và các công cụ quản lý DAG trực quan của IDE.
Khái niệm chính
Khái niệm | Kiến thức bạn học được |
Lập trình theo cặp trong IDE bằng ngôn ngữ tự nhiên để tạo sổ tay PySpark, định cấu hình các mô hình dbt và xác định DAG Airflow | |
Bộ nhớ dạng bảng có thể mở rộng cho SQL phân tích, các phép biến đổi dbt và hoạt động huấn luyện mô hình học máy | |
Thực thi không máy chủ để tải dữ liệu PySpark phân tán và huấn luyện ML Random Forest | |
Viết các dự đoán suy luận hàng loạt của Spark trực tiếp vào hàng đợi xem xét cơ sở dữ liệu hoạt động | |
Khai báo DAG YAML | Các định nghĩa về quy trình khai báo được hiển thị dưới dạng biểu đồ trực quan tương tác của Airflow trong IDE |
Quản lý DAG trực quan | Kiểm tra các phần phụ thuộc của pipeline, triển khai đến Managed Airflow và giám sát nhật ký thực thi tác vụ trực tiếp trong IDE |
Các bước tiếp theo
- Khám phá tài liệu về Bộ công cụ tác nhân dữ liệu của Google Cloud
- Tìm hiểu thêm về Dịch vụ được quản lý cho Apache Spark
- Tìm hiểu thêm về Dịch vụ được quản lý cho Apache Airflow
- Tạo các quy trình đa dịch vụ của riêng bạn bằng Antigravity IDE