บทนำสู่ Dask DataFrames
แทนที่ pd.read_csv และ pd.DataFrame ด้วยสิ่งเทียบเท่าใน dask เรียกใช้ compute() เพื่อเริ่มการประมวลผล และวิเคราะห์กราฟงาน
บทนำสู่ Dask DataFrames เป็นบทเรียน Pandas & NumPy Academy ฟรีบน CoddyKit นี่คือบทเรียนที่ 3 จากทั้งหมด 4 บทเรียน คุณสามารถอ่านบทเรียนทั้งหมดด้านล่างฟรี — จากนั้นลองปฏิบัติด้วยตัวคุณเองในเบราว์เซอร์พร้อมตัวแก้ไขโค้ดในตัวและติวเตอร์ AI ตลอด 24/7 บทเรียนนี้เป็นส่วนหนึ่งของเส้นทางการเรียน Pandas & NumPy Academy และความก้าวหน้าของคุณจะซิงค์ข้ามเว็บและแอป CoddyKit คอร์ส Pandas & NumPy Academy มีบทเรียนทั้งหมด 4 บทเรียน
Dask คืออะไร
Dask เป็นไลบรารีการประมวลผลแบบขนานสำหรับ Python ที่ขยายความสามารถของ NumPy และ Pandas ให้รองรับชุดข้อมูลที่มีขนาดใหญ่กว่า RAM โมดูล dask.dataframe มีส่วนติดต่อการใช้งาน DataFrame ที่เกือบเหมือนกับ Pandas ทุกประการ แต่แทนที่จะดำเนินการทันที Dask จะสร้างกราฟงานและดำเนินการแบบเลื่อนเวลาไว้จนกว่าคุณจะเรียกใช้ .compute() ทำให้ Dask สามารถแบ่งงานไปทำงานบนหลายแกนประมวลผลหรือแม้แต่หลายเครื่องได้ โดยเปลี่ยนแปลงโค้ดเพียงเล็กน้อย
การติดตั้งและนำเข้า Dask
ติดตั้ง Dask ด้วย pip install dask[dataframe] รูปแบบการนำเข้าที่ใช้ทั่วไปคือ import dask.dataframe as dd เบื้องหลังนั้น Dask DataFrame จะถูกแบ่งเป็น Pandas DataFrame ขนาดเล็กหลายชุด โดยแต่ละชุดจะถูกประมวลผลแยกจากกัน การดำเนินการกับ Dask DataFrame จะสร้างกราฟงานแบบเลื่อนเวลา — จะยังไม่มีอะไรทำงานจนกว่าจะเรียกใช้ .compute() การแยกการอธิบายการคำนวณออกจากการดำเนินการคำนวณนี้คือแนวคิดสำคัญของ Dask
import dask.dataframe as dd
# Read a large CSV — returns a Dask DataFrame immediately (no data loaded yet)
ddf = dd.read_csv('large_sales.csv')
print(type(ddf)) # dask.dataframe.core.DataFrame
print(ddf.columns.tolist())
print(ddf.dtypes)Dask กับ Pandas: ความแตกต่างสำคัญ
สำหรับ Pandas การดำเนินการทุกอย่างจะทำงานทันทีและโดยอัตโนมัติ ส่วน Dask จะส่งคืนออบเจ็กต์ Dask อีกตัวที่แสดงการคำนวณซึ่งเลื่อนเวลาไว้ Dask จะอ่านข้อมูลและดำเนินการตามกราฟงานจริงก็ต่อเมื่อคุณเรียกใช้ .compute() ความล่าช้านี้ทำให้ Dask สามารถปรับแผนให้เหมาะสมก่อนดำเนินการได้ ตัวอย่างเช่น สามารถรวมตัวกรองที่ต่อเนื่องกันเพื่อลดการโหลดข้อมูลซ้ำหลายครั้ง ลองนึกภาพว่าเป็นสูตรอาหาร: Dask เขียนสูตร ส่วน .compute() เป็นผู้ปรุงอาหาร
import dask.dataframe as dd
ddf = dd.read_csv('sales.csv')
# This does NOT run yet — just builds the task graph
filtered = ddf[ddf['amount'] > 1000]
agg = filtered.groupby('region')['amount'].sum()
print(type(agg)) # dask.dataframe.core.Series
# NOW execute everything
result = agg.compute()
print(result)พาร์ติชัน: แนวคิดหลัก
Dask DataFrame ถูกแบ่งเป็นพาร์ติชัน โดยแต่ละพาร์ติชันคือ Pandas DataFrame ปกติ โดยค่าเริ่มต้น dd.read_csv จะสร้างหนึ่งพาร์ติชันต่อหนึ่งไฟล์ (หรือหนึ่งพาร์ติชันต่อ 128 MB สำหรับไฟล์ขนาดใหญ่) คุณสามารถควบคุมค่านี้ได้ด้วย blocksize การตรวจสอบ ddf.npartitions จะแสดงจำนวนพาร์ติชันที่มีอยู่ พาร์ติชันจำนวนมากช่วยให้ทำงานแบบขนานได้มากขึ้น แต่เพิ่มภาระการจัดการ ส่วนพาร์ติชันจำนวนน้อยจะลดภาระดังกล่าว แต่จำกัดการทำงานแบบขนาน โดยทั่วไปจำนวนที่เหมาะสมมักอยู่ที่ไม่กี่ร้อยพาร์ติชัน
import dask.dataframe as dd
ddf = dd.read_csv('data/*.csv') # Read multiple CSV files at once
print('Number of partitions:', ddf.npartitions)
# Access a single partition as a Pandas DataFrame
first_partition = ddf.get_partition(0).compute()
print('Partition 0 shape:', first_partition.shape)การดำเนินการของ Pandas ที่คุ้นเคยใน Dask
การดำเนินการทั่วไปของ Pandas ส่วนใหญ่ทำงานใน Dask ได้เหมือนกัน: .head(), .tail(), .describe(), การเลือกข้อมูลด้วยบูลีน, .groupby(), .merge() และ .assign() ล้วนมีสิ่งเทียบเท่าใน Dask ความแตกต่างที่สำคัญที่สุดคือ คุณต้องเรียกใช้ .compute() เพื่อทำให้ผลลัพธ์เกิดขึ้นจริง การดำเนินการที่ Pandas ใช้เวลาเพียงไม่กี่มิลลิวินาทีอาจใช้เวลาหลายวินาทีใน Dask เนื่องจากภาระของกราฟงาน ดังนั้นควรใช้ Pandas กับข้อมูลขนาดเล็ก และใช้ Dask เมื่อข้อมูลมีขนาดใหญ่เกิน RAM
import dask.dataframe as dd
ddf = dd.read_csv('transactions.csv')
# Filtering — same syntax as Pandas
high_value = ddf[ddf['amount'] > 500]
# GroupBy aggregation
by_region = high_value.groupby('region')['amount'].mean()
# Execute
result = by_region.compute()
print(result.sort_values(ascending=False))การอ่านหลายไฟล์ด้วยรูปแบบ Glob
หนึ่งในคุณสมบัติที่มีประโยชน์ที่สุดของ Dask คือการอ่านหลายไฟล์พร้อมกันโดยใช้รูปแบบ glob dd.read_csv('data/2024-*.csv') จะอ่านไฟล์ทั้งหมดที่ตรงกับรูปแบบและสร้างหนึ่งพาร์ติชันต่อหนึ่งไฟล์ วิธีนี้เหมาะอย่างยิ่งกับข้อมูลที่จัดเก็บเป็นไฟล์แยกตามเดือนหรือวัน ซึ่งเป็นรูปแบบที่พบได้ทั่วไปในคลังข้อมูล Dask จะจัดแนวโครงสร้างข้อมูลให้ตรงกันโดยอัตโนมัติ เทียบเท่ากับการวนซ้ำและต่อข้อมูลด้วย Pandas เอง แต่ทำได้ง่ายกว่ามาก
import dask.dataframe as dd
# Read all monthly files at once
ddf = dd.read_csv('sales/2024-*.csv',
dtype={'order_id': 'int32',
'amount': 'float32'})
print(f'Partitions: {ddf.npartitions}') # One per file
print(f'Total rows (lazy): {len(ddf)}') # This triggers a compute!เมธอด visualize() สำหรับกราฟงาน
ก่อนดำเนินการกับกระบวนการข้อมูลของ Dask ที่ซับซ้อน คุณสามารถตรวจสอบกราฟงานได้โดยเรียกใช้ result.visualize() ซึ่งจะสร้างแผนภาพ PNG ที่แสดงขั้นตอนการคำนวณทั้งหมด วิธีนี้มีประโยชน์สำหรับทำความเข้าใจว่า Dask จะดำเนินการอะไรบ้าง และช่วยแก้ไขปัญหาความล่าช้าที่ไม่คาดคิด กราฟจะแสดงการไหลของพาร์ติชันผ่านขั้นตอนการกรอง, groupby และการรวมข้อมูล ทำให้มองเห็นการคำนวณซ้ำซ้อนได้ง่าย ต้องติดตั้งแพ็กเกจ graphviz
import dask.dataframe as dd
ddf = dd.read_csv('orders.csv')
pipeline = (
ddf[ddf['status'] == 'completed']
.groupby('product_id')['revenue']
.sum()
)
# Visualise the task graph (saves to PNG)
# pipeline.visualize('task_graph.png')
# Check number of tasks in the graph
print('Number of tasks:', len(pipeline.__dask_graph__()))การใช้ฟังก์ชันกำหนดเองด้วย map_partitions
เมื่อคุณต้องการใช้ฟังก์ชัน Pandas ที่กำหนดเองกับ Dask DataFrame ให้ใช้ ddf.map_partitions(func) การดำเนินการนี้จะใช้ func กับแต่ละพาร์ติชันแยกจากกัน และส่งคืน Dask DataFrame ตัวใหม่ ฟังก์ชันจะได้รับ Pandas DataFrame ปกติและต้องส่งคืน DataFrame เช่นกัน วิธีนี้เป็นสิ่งเทียบเท่าของ Dask กับ df.apply() และเป็นวิธีผสาน Dask เข้ากับโค้ดที่เข้าใจเฉพาะ Pandas
import dask.dataframe as dd
import pandas as pd
def normalise_chunk(df):
df = df.copy()
df['amount_norm'] = (df['amount'] - df['amount'].mean()) / df['amount'].std()
return df
ddf = dd.read_csv('data.csv')
normalised = ddf.map_partitions(normalise_chunk)
result = normalised[['order_id', 'amount_norm']].compute()
print(result.head())ตัวเลือกตัวจัดตารางงานของ Dask
Dask มีตัวจัดตารางงานหลายแบบสำหรับควบคุมวิธีดำเนินการของงาน ตัวจัดตารางงาน 'synchronous' จะทำงานตามลำดับในเธรดปัจจุบัน (มีประโยชน์สำหรับการแก้ไขข้อผิดพลาด) ตัวจัดตารางงาน 'threads' ใช้กลุ่มเธรด (เหมาะกับงานที่รอการอ่านเขียนข้อมูล) ตัวจัดตารางงาน 'processes' สร้างกระบวนการหลายกระบวนการสำหรับงานที่ใช้ CPU เป็นหลัก (หลีกเลี่ยง GIL ของ Python) คลัสเตอร์ Dask แบบกระจายช่วยให้ทำงานบนหลายเครื่องได้ ระบุตัวจัดตารางงานผ่าน compute(scheduler='threads')
import dask.dataframe as dd
ddf = dd.read_csv('data.csv')
agg = ddf.groupby('category')['sales'].sum()
# Choose scheduler based on workload
result_sync = agg.compute(scheduler='synchronous') # sequential, easy to debug
result_threads = agg.compute(scheduler='threads') # parallel I/O
result_processes = agg.compute(scheduler='processes') # parallel CPUการแปลงระหว่าง Dask และ Pandas
เป็นเรื่องปกติที่จะประมวลผลชุดข้อมูลขนาดใหญ่ด้วย Dask แล้วนำผลลัพธ์ที่รวมแล้วเข้าสู่ Pandas เพื่อวิเคราะห์หรือแสดงผลในขั้นสุดท้าย ใช้ .compute() เพื่อแปลง Dask DataFrame เป็น Pandas ในทางกลับกัน dd.from_pandas(df, npartitions=4) จะแปลง Pandas DataFrame เป็น Dask DataFrame ซึ่งมีประโยชน์สำหรับทดสอบโค้ด Dask กับข้อมูลขนาดเล็กก่อนขยายไปใช้กับชุดข้อมูลทั้งหมด
import pandas as pd
import dask.dataframe as dd
# Start with a small Pandas DF for testing
df_small = pd.DataFrame({'a': range(100), 'b': range(100, 200)})
# Convert to Dask for development/testing
ddf = dd.from_pandas(df_small, npartitions=4)
result = ddf.groupby('a')['b'].sum().compute()
print(type(result)) # pandas.Series
print(result.head())ควรใช้ Dask, Pandas หรือ SQL เมื่อใด
Dask ไม่ใช่เครื่องมือที่เหมาะสมเสมอไป ใช้ Pandas เมื่อข้อมูลของคุณมีขนาดพอดีกับ RAM (น้อยกว่าไม่กี่ GB) — ใช้งานง่ายกว่าและเร็วกว่าเนื่องจากมีภาระการจัดการน้อยกว่า ใช้ Dask เมื่อข้อมูลมีขนาดใหญ่เกิน RAM แต่คุณต้องการรูปแบบไวยากรณ์คล้าย Pandas และการประมวลผลแบบขนานบนเครื่องเดียว ใช้ SQL/ฐานข้อมูล เมื่อข้อมูลอยู่ในฐานข้อมูลเชิงสัมพันธ์และสามารถส่งการรวมข้อมูลไปให้กลไกฐานข้อมูลดำเนินการได้ ใช้ Spark หรือ BigQuery เมื่อต้องการการประมวลผลแบบกระจายหลายเครื่องกับข้อมูลขนาดระดับเพตะไบต์
ตรวจสอบความเข้าใจอย่างรวดเร็ว
ทดสอบความเข้าใจแนวคิดการวิเคราะห์ข้อมูลจากบทเรียนนี้
ทบทวนบทเรียน
ในบทเรียนนี้ คุณได้เรียนรู้ว่า Dask DataFrames คือคอลเลกชันของพาร์ติชัน Pandas ที่ประเมินผลแบบเลื่อนเวลาไว้และมีส่วนติดต่อการใช้งานที่คุ้นเคย, .compute() จะเริ่มการดำเนินการจริงของกราฟงาน และ map_partitions ช่วยให้คุณใช้ฟังก์ชัน Pandas ที่กำหนดเองกับทุกพาร์ติชันได้ ต่อไปเราจะดูรูปแบบ Parquet ซึ่งเป็นทางเลือกแบบคอลัมน์ที่รวดเร็วกว่า CSV สำหรับจัดเก็บชุดข้อมูลขนาดใหญ่
คำถามที่พบบ่อย
บทเรียน “บทนำสู่ Dask DataFrames” ฟรีหรือไม่
ใช่ — ข้อความเต็มของ “บทนำสู่ Dask DataFrames” ฟรีให้อ่านที่นี่บนเว็บ เพื่อปฏิบัติแบบโต้ตอบ (ตัวแก้ไขโค้ดในตัวและติวเตอร์ AI ตลอด 24/7) และปลดล็อคส่วนที่เหลือของคอร์ส Pandas & NumPy Academy ให้อัปเกรดเป็น CoddyKit PRO คอร์ส Pandas & NumPy Academy มีบทเรียนทั้งหมด 4 บทเรียน
คุณจะเรียนรู้อะไรในบทเรียน “บทนำสู่ Dask DataFrames”
แทนที่ pd.read_csv และ pd.DataFrame ด้วยสิ่งเทียบเท่าใน dask เรียกใช้ compute() เพื่อเริ่มการประมวลผล และวิเคราะห์กราฟงาน คุณปฏิบัติ Pandas & NumPy Academy ด้วยโค้ดที่ใช้งานได้จริงที่คุณเรียกใช้โดยตรงในเบราว์เซอร์ และติวเตอร์ AI ตลอด 24/7 ตอบคำถามของคุณขณะที่คุณไปผ่านบทเรียน
คุณต้องมีประสบการณ์ก่อนที่จะเริ่มเรียน Pandas & NumPy Academy หรือไม่
ไม่จำเป็นต้องมีประสบการณ์มาก่อน Pandas & NumPy Academy บน CoddyKit ออกแบบมาสำหรับผู้เริ่มต้นไปจนถึงผู้เรียนขั้นสูง คุณสามารถเริ่มต้นที่นี่หรือเริ่มจากตัวแรกและเรียนด้วยความเร็วของคุณเอง นี่คือบทเรียนที่ 3 จากทั้งหมด 4 บทเรียน
บทเรียน “บทนำสู่ Dask DataFrames” ใช้เวลานานแค่ไหน
บทเรียน CoddyKit ส่วนใหญ่ใช้เวลาประมาณ 5–10 นาที แต่ละบทเรียนจึงสั้นและเป็นแบบโต้ตอบ คุณสามารถก้าวหน้าอย่างต่อเนื่องและกลับมาเรียนต่อจากตรงที่เพิ่งหยุดบนเว็บและแอปได้เลย
ฉันเขียนและรันโค้ดในบทเรียน Pandas & NumPy Academy นี้ได้ไหม
ได้ บทเรียน Pandas & NumPy Academy ทุกบทมีตัวแก้ไขโค้ดในตัว คุณจึงเขียนและรันโค้ดจริงได้เลยในเบราว์เซอร์ และได้รับข้อเสนอแนะจาก AI ในทันที — ไม่ต้องติดตั้งในเครื่องของคุณ
บทเรียนทั้งหมดในหลักสูตรนี้
- การสตรีม CSV ด้วย chunksize
- การรวมค่าแบบเพิ่มทีละส่วนระหว่างชิ้นข้อมูล
- บทนำสู่ Dask DataFrames
- Parquet: การจัดเก็บข้อมูลแบบคอลัมน์ความเร็วสูง