ML trên Snowflake ở quy mô lớn với Snowpark Python và XGBoost
Trong hơn một năm rưỡi qua, nhóm của chúng tôi đã làm việc với hàng trăm khách hàng xây dựng trên Snowflake với Snowpark và Python (nhân tiện - hãy xem bài đăng này của đồng nghiệp Caleb Baechtold của tôi , tổng hợp tất cả các bài học về cách vận hành Snowpark Python trong sản xuất) và một câu hỏi tôi thường nhận được là, "nhưng chúng tôi có dữ liệu thực sự lớn, nó hoạt động như thế nào trên quy mô lớn?".
Nhiều ví dụ/bản trình diễn mà bạn sẽ tìm thấy trực tuyến rất tuyệt khi hiển thị các ví dụ giả tạo đơn giản thể hiện chức năng cốt lõi, nhưng rất ít nếu có thể hiện các ví dụ về cách thức hoạt động của nó trong bối cảnh các vấn đề ở quy mô doanh nghiệp lớn.
TPC-DS , là một tập dữ liệu quy mô lớn hữu ích thường đại diện cho một lược đồ dữ liệu kinh doanh chung (Tôi không ở đây để ủng hộ rằng TPC-DS hoặc bất kỳ tập dữ liệu nào là mục tiêu cuối cùng của bất kỳ nền tảng/hoặc công cụ nào, tôi chỉ nghĩ rằng nó có thể là một tập dữ liệu hữu ích làm điểm khởi đầu để kiểm tra mọi thứ ở quy mô lớn). Snowflake cung cấp TPC-DS trong mọi tài khoản Snowflake dưới dạng chia sẻ dữ liệu, ở cả phiên bản 10 TB và 100 TB (nhân tiện, vì nó được hiển thị dưới dạng chia sẻ dữ liệu mà bạn là người dùng không phải trả bất kỳ chi phí nào cho việc lưu trữ . Chỉ dành cho mọi điện toán tiềm năng mà bạn sử dụng để truy vấn nó). Phiên bản 100 TB có hơn 560 tỷ hàng trong bảng thực tế. Bản 10 TB, 56 tỷ.
Chúng ta có thể sử dụng Snowpark Python để xây dựng một giải pháp ML tương đối đơn giản cho một vấn đề kinh doanh phổ biến mà một doanh nghiệp giả định như TPC gặp phải, đó là “Tôi muốn có thể dự đoán giá trị lâu dài của khách hàng trên tất cả các kênh bán hàng”. Chúng tôi sẽ sử dụng API Snowpark Python Dataframe để thực hiện kỹ thuật tính năng/chuẩn bị dữ liệu, quy trình được lưu trữ và Kho được tối ưu hóa của Snowpark để đào tạo và xử lý hàng loạt UDF để suy luận mà không cần dữ liệu rời khỏi Snowflake bằng cách sử dụng tài nguyên điện toán và khả năng mở rộng quy mô của Snowflake.
Hãy bắt đầu với mã kỹ thuật tính năng/chuẩn bị dữ liệu của chúng tôi. Khá đơn giản, chúng tôi sẽ tổng hợp doanh số bán hàng theo khách hàng trên tất cả các kênh. Sau đó, chúng tôi sẽ kết hợp điều đó với bảng thứ nguyên khách hàng của chúng tôi để có được các tính năng tiềm năng được quan tâm.
store_sales_agged = store_sales.group_by('ss_customer_sk').agg(F.sum('ss_sales_price').as_('total_sales'))
web_sales_agged = web_sales.group_by('ws_bill_customer_sk').agg(F.sum('ws_sales_price').as_('total_sales'))
catalog_sales_agged = catalog_sales.group_by('cs_bill_customer_sk').agg(F.sum('cs_sales_price').as_('total_sales'))
store_sales_agged = store_sales_agged.rename('ss_customer_sk', 'customer_sk')
web_sales_agged = web_sales_agged.rename('ws_bill_customer_sk', 'customer_sk')
catalog_sales_agged = catalog_sales_agged.rename('cs_bill_customer_sk', 'customer_sk')
total_sales = store_sales_agged.union_all(web_sales_agged)
total_sales = total_sales.union_all(catalog_sales_agged)
total_sales = total_sales.group_by('customer_sk').agg(F.sum('total_sales').as_('total_sales'))
customer = customer.select('c_customer_sk','c_current_hdemo_sk', 'c_current_addr_sk', 'c_customer_id', 'c_birth_year')
customer = customer.join(address.select('ca_address_sk', 'ca_zip'), customer['c_current_addr_sk'] == address['ca_address_sk'] )
customer = customer.join(demo.select('cd_demo_sk', 'cd_gender', 'cd_marital_status', 'cd_credit_rating', 'cd_education_status', 'cd_dep_count'),
customer['c_current_hdemo_sk'] == demo['cd_demo_sk'] )
customer = customer.rename('c_customer_sk', 'customer_sk')
final_df = total_sales.join(customer, on='customer_sk')
session.use_database('tpcds_xgboost')
session.use_schema('demo')
final_df.write.mode('overwrite').save_as_table('feature_store')
Bây giờ chúng tôi đã sẵn sàng để đào tạo mô hình của mình bằng cách sử dụng thủ tục được lưu trữ. Snowflake chỉ đơn giản là cung cấp tất cả các khung Python ML phổ biến để bạn sử dụng cho đào tạo ML thông qua quan hệ đối tác của chúng tôi với Anaconda . Không cần phải học một thư viện hoặc cú pháp mới.
from sklearn.pipeline import Pipeline
from sklearn.impute import SimpleImputer
from sklearn.preprocessing import StandardScaler, OneHotEncoder, MinMaxScaler
from sklearn.metrics import mean_squared_error
from sklearn.compose import ColumnTransformer
from xgboost import XGBRegressor
import joblib
import os
def train_model(session: snowflake.snowpark.Session) -> float:
snowdf = session.table("feature_store")
snowdf = snowdf.drop(['CUSTOMER_SK', 'C_CURRENT_HDEMO_SK', 'C_CURRENT_ADDR_SK', 'C_CUSTOMER_ID', 'CA_ADDRESS_SK', 'CD_DEMO_SK'])
snowdf_train, snowdf_test = snowdf.random_split([0.8, 0.2], seed=82)
# save the train and test sets as time stamped tables in Snowflake
snowdf_train.write.mode("overwrite").save_as_table("tpcds_xgboost.demo.tpc_TRAIN")
snowdf_test.write.mode("overwrite").save_as_table("tpcds_xgboost.demo.tpc_TEST")
train_x = snowdf_train.drop("TOTAL_SALES").to_pandas() # drop labels for training set
train_y = snowdf_train.select("TOTAL_SALES").to_pandas()
test_x = snowdf_test.drop("TOTAL_SALES").to_pandas()
test_y = snowdf_test.select("TOTAL_SALES").to_pandas()
cat_cols = ['CA_ZIP', 'CD_GENDER', 'CD_MARITAL_STATUS', 'CD_CREDIT_RATING', 'CD_EDUCATION_STATUS']
num_cols = ['C_BIRTH_YEAR', 'CD_DEP_COUNT']
num_pipeline = Pipeline([
('imputer', SimpleImputer(strategy="median")),
('std_scaler', StandardScaler()),
])
preprocessor = ColumnTransformer(
transformers=[('num', num_pipeline, num_cols),
('encoder', OneHotEncoder(handle_unknown="ignore"), cat_cols) ])
pipe = Pipeline([('preprocessor', preprocessor),
('xgboost', XGBRegressor())])
pipe.fit(train_x, train_y)
test_preds = pipe.predict(test_x)
rmse = mean_squared_error(test_y, test_preds)
model_file = os.path.join('/tmp', 'model.joblib')
joblib.dump(pipe, model_file)
session.file.put(model_file, "@ml_models",overwrite=True)
return rmse
train_model_sp = F.sproc(train_model, session=session, replace=True)
# Switch to Snowpark Optimized Warehouse for training and to run the stored proc
session.use_warehouse('snowpark_opt_wh')
train_model_sp(session=session)
import sys
import pandas as pd
import cachetools
import joblib
from snowflake.snowpark import types as T
session.add_import("@ml_models/model.joblib")
features = [ 'C_BIRTH_YEAR', 'CA_ZIP', 'CD_GENDER', 'CD_MARITAL_STATUS', 'CD_CREDIT_RATING', 'CD_EDUCATION_STATUS', 'CD_DEP_COUNT']
@cachetools.cached(cache={})
def read_file(filename):
import_dir = sys._xoptions.get("snowflake_import_directory")
if import_dir:
with open(os.path.join(import_dir, filename), 'rb') as file:
m = joblib.load(file)
return m
@F.pandas_udf(session=session, max_batch_size=10000, is_permanent=True,
stage_location='@ml_models', name="clv_xgboost_udf")
def predict(df: T.PandasDataFrame[int, str, str, str, str, str, int]) -> T.PandasSeries[float]:
m = read_file('model.joblib')
df.columns = features
return m.predict(df)
inference_df = session.table('feature_store')
inference_df = inference_df.drop(['CUSTOMER_SK', 'C_CURRENT_HDEMO_SK', 'C_CURRENT_ADDR_SK', 'C_CUSTOMER_ID', 'CA_ADDRESS_SK', 'CD_DEMO_SK'])
inputs = inference_df.drop("TOTAL_SALES")
snowdf_results = inference_df.select(*inputs,
predict(*inputs).alias('PREDICTION'),
(F.col('TOTAL_SALES')).alias('ACTUAL_SALES')
)
snowdf_results.write.mode('overwrite').save_as_table('predictions')
SELECT "C_BIRTH_YEAR",
"CA_ZIP",
"CD_GENDER",
"CD_MARITAL_STATUS",
"CD_CREDIT_RATING",
"CD_EDUCATION_STATUS",
"CD_DEP_COUNT",
clv_xgboost_udf("C_BIRTH_YEAR", "CA_ZIP", "CD_GENDER", "CD_MARITAL_STATUS", "CD_CREDIT_RATING", "CD_EDUCATION_STATUS", "CD_DEP_COUNT") AS "PREDICTION",
"TOTAL_SALES" AS "ACTUAL_SALES"
FROM tpcds_xgboost.demo.feature_store
Khi tôi nói với tất cả các khách hàng mà tôi làm việc cùng, xin đừng tin lời tôi nói. Hãy tự mình thử, tất cả mã đều có sẵn cho bạn tại đây . Thậm chí còn tốt hơn TPC-DS, hãy dùng thử với một số dữ liệu và quy trình của tổ chức bạn.

![Dù sao thì một danh sách được liên kết là gì? [Phần 1]](https://post.nghiatu.com/assets/images/m/max/724/1*Xokk6XOjWyIGCBujkJsCzQ.jpeg)



































