ورودی بزرگA large input
کار یک فایل، یک پوشه، یا یک جدول کامل را دریافت میکند؛ نه یک شناسه در URL.The job takes a file, a folder, or a whole table. Not one id in a URL.
داده را یکجا و در قالب یک کار پردازش کنید — نه رکورد به رکورد و آنلاین. API شما همیشه پاسخگوست، اما گزارش شبانه در پسزمینه میچرخد. این دو دنیای کاملاً جدا هستند. Process data in bulk as one job — not record by record in real time. Your API is online. The nightly report is offline. These are two different worlds.
پردازش دستهای یعنی داده را بهصورت انبوه و یکجا پردازش کنید؛ نه رکورد به رکورد و لحظهای. کتاب Designing Data-Intensive Applications این رویکرد را در برابر سیستمهای آنلاین قرار میدهد: در سیستم آنلاین، درخواستی میرسد و بلافاصله پاسخ میگیرد؛ در پردازش آفلاین (دستهای)، یک کار، تمام ورودی را میخواند و در پایان خروجی را مینویسد. Batch processing means you process data in bulk as one job. You do not handle each record in real time. Designing Data-Intensive Applications contrasts this with online systems: online is request and response. Offline is a job that reads all input and writes an output.
در معماری داده، این کار یکی از قطعههای یک پازل بزرگتر است: معماریِ خوب اجزای پرکاربرد را هوشمندانه برمیگزیند، برای خرابی برنامهریزی میکند و سیستمها را از هم مستقل نگه میدارد. کار دستهای معمولاً همان قطعهای است که دادهٔ عملیاتی را به دادهٔ تحلیلی تبدیل میکند؛ بیآنکه بار اضافهای روی API فروش بگذارد. In a data architecture, this job is one piece of a larger picture. Good architecture picks common parts wisely, plans for failure, and keeps systems loosely coupled. A batch job is often the piece that turns operational data into analytical data — without loading down the sales API.
کاربر سفارش میدهد. API در کمتر از یک ثانیه پاسخ میدهد. هر رکورد در همین لحظه اهمیت دارد.A user places an order. The API answers in under a second. Each record matters now.
ساعت ۲ بامداد، همهٔ سفارشهای دیروز خوانده میشوند و یک گزارش روزانه ساخته میشود.At 2 a.m., all of yesterday's orders are read. One daily report is built.
کار یک فایل، یک پوشه، یا یک جدول کامل را دریافت میکند؛ نه یک شناسه در URL.The job takes a file, a folder, or a whole table. Not one id in a URL.
اگر کار ۲۰ دقیقه طول بکشد، کسی صفحه را تازهسازی نمیکند. زمان کل اهمیت دارد، نه تأخیر یک رکورد.If the job takes 20 minutes, nobody refreshes a page. Total time matters, not per-record latency.
همان ورودی باید همان خروجی را تولید کند. اگر کار خراب شد، دوباره اجرایش میکنید.The same input should give the same output. If the job fails, you run it again.
فروشگاه شما روزی ۴۰۰٬۰۰۰ سفارش ثبت میکند. API هر سفارش را در ۱۰ میلیثانیه مینویسد — این همان پردازش آنلاین است. گزارش شبانه، همان ۴۰۰٬۰۰۰ سطر را یکجا جمع میزند و در ۸ دقیقه تمام میشود. اگر همان جمع را داخل هر درخواست API محاسبه کنید، تسویهٔ سبد خرید با کندی مواجه میشود.Your shop takes 400,000 orders a day. The API writes each order in 10 ms — that is online. The nightly report sums those 400,000 rows in one job and finishes in 8 minutes. If you compute that sum inside every API call, checkout gets slow.
برای یک توسعهدهندهٔ بکاند داتنت این یعنی: BackgroundService، Hangfire، یا یک Azure Function زمانبندیشده. کنترلر MVC برای این کار مناسب نیست. درخواست HTTP باید کوتاه بماند. کار دستهای را جداگانه اجرا کنید و نتیجه را در جدول گزارش یا فایل بنویسید.For a .NET backend dev this means: a BackgroundService, Hangfire, or a scheduled Azure Function. An MVC controller is the wrong place. Keep the HTTP request short. Run the batch job separately and write the result to a report table or a file.
نیمهٔ دوم فصل ۴ کتاب DDIA به یک نکتهٔ مهم اشاره میکند: کارهای تحلیلی به چیدمان دیگری برای دادهها روی دیسک نیاز دارند. سیستمهای آنلاین (OLTP) هر سطر را کامل میخوانند؛ مثلاً یک سفارش با همهٔ ستونهایش. اما کار تحلیلی معمولاً فقط چند ستون خاص را از میان میلیونها سطر جمع میکند. برای چنین کاری، ذخیرهسازی ستونمحور بهمراتب بهصرفهتر است. The second half of DDIA Chapter 4 says analytical work needs a different layout on disk. An online system (OLTP) wants a whole row: one order, every column. An analytical job usually sums a few columns across millions of rows. For that, column-oriented storage is better.
پیشتر، در صفحهٔ ایندکسها با B-tree و LSM-tree و ایندکس هش آشنا شدهاید؛ اینجا دوباره به آنها نمیپردازیم. نکتهٔ تازه این است: موتوری که برای یافتن سریع یک سفارش عالی است — همان ذخیرهسازی سطرمحور — برای محاسبهٔ «جمع فروش کل سال» بسیار پرهزینه میشود. You already met B-trees, LSM-trees, and hash indexes on the Indexes page. We will not teach them again. The new point: the same row store that finds one order fast is expensive for "sum of sales this year."
سطرهای جدول کنار هم روی دیسک مینشینند؛ خواندن یک سفارش کامل ارزان است، اما خواندن یک ستون از ۱۰ میلیون سطر عملاً یعنی عبور از بیشترِ جدول.Each row sits together on disk. Reading one order is cheap. Reading one column from 10 million rows means you touch most of the table.
هر ستون در فایل خودش ذخیره میشود. موتور فقط ستون Amount را میخواند. فشردهسازی هم بهتر کار میکند، چون مقادیر مشابه کنار هم قرار میگیرند.Each column lives in its own file. The engine reads only Amount. Compression also works better, because similar values sit together.
جدول کوچکِ پایین ۴ ستون و ۶ سطر دارد و کوئری فقط جمع Amount را میخواهد. سطرمحور تقریباً همهٔ سلولها را میخواند، اما ستونمحور فقط یک ستون را. یکی از دو حالت را انتخاب کنید تا اسکن اجرا شود.The tiny table below has 4 columns and 6 rows. The query only wants SUM(Amount). A row store reads almost everything. A column store reads one column. Pick a layout to run the scan.
کار تحلیلی معمولاً روی انبار داده یا موتور ستونی اجرا میشود، نه روی SQL Server عملیاتی شما. فصل ۴ همچنین به کامپایل کوئری، پردازش برداری (Vectorization)، نمای مادی (Materialized View) و مکعب داده (Data Cube) اشاره میکند. ایده یکی است: هزینهٔ یک اسکن بزرگ را پایین بیاورید، نه هزینهٔ یافتن یک سطر.Analytical jobs often run on a warehouse or a columnar engine, not on your operational SQL Server. Chapter 4 also covers query compilation, vectorization, materialized views, and data cubes. The idea is the same: make a large scan cheap, not a single-row lookup.
اگر گزارش سنگین را مستقیم روی دیتابیس سفارشها اجرا کنید، قفلها و I/O با API تداخل پیدا میکنند. الگوی رایج این است: شب، داده را به یک انبار (مثل Synapse، BigQuery، یا حتی یک دیتابیس فقطخواندنی ستونی) کپی کنید و برنامهٔ داتنت گزارش را از آنجا بخواند.If you run a heavy report on the order database, locks and IO mix with the API. A common pattern: copy data at night into a warehouse (Synapse, BigQuery, or even a read-only columnar store). The .NET app reads reports from there.
بعد از تعریف کار دستهای و محل ذخیرهٔ دادهٔ تحلیلی، حالا نوبت خودِ پردازش است؛ و سادهترین شکل آن، خطلولهٔ یونیکس است. فصل ۱۱ کتاب DDIA با ابزارهای خط فرمان شروع میشود: لاگ را با grep، awk، sort و uniq به هم زنجیر میکنید؛ هر برنامه یک کار کوچک انجام میدهد و خروجی هر کدام، ورودی مرحلهٔ بعد میشود. این همان مدل ذهنی پردازش دستهای است.
Now that we have defined the batch job and where analytical data lives, it is time for the processing itself — and its simplest form is the Unix pipeline. Chapter 11 starts with Unix tools. You chain a log through grep, awk, sort, and uniq. Each program does one small job, and one output becomes the next input. That is the mental model of a batch job.
# How many paid orders per country yesterday?
cat orders.log \
| grep 'status=paid' \
| awk -F, '{print $3}' \
| sort \
| uniq -c
هر مرحله را میتوانید جداگانه تست و جایگزین کنید؛ اگر sort کند بود، فقط همان مرحله را عوض میکنید. کتاب، این رویکرد را به برنامهٔ بزرگ و یکپارچهای که همهچیز را یکجا انجام میدهد ترجیح میدهد.You test each step alone. If sort is slow, you replace only that step. The book sets this against one custom program.
اگر همهٔ کلیدها در RAM جا شوند، یک دیکشنری کافی است. اگر جا نشوند، sort داده را روی دیسک میریزد و سپس گروهها را میسازد. سیستمهای بزرگ هم دقیقاً همین کار را میکنند.If every key fits in RAM, a dictionary is enough. If not, sort spills to disk and then builds groups. Large jobs do the same thing.
وقتی حجم داده از یک ماشین فراتر میرود، دیگر نمیتوانید همهچیز را روی یک دیسک و با یک پردازنده انجام دهید. باید داده را بین چند ماشین پخش کنید. برای این کار، سه رویکرد اصلی وجود دارد که در جدول زیر مقایسه شدهاند. انتخاب هر کدام به حجم داده، بودجه و نیاز شما به سرعت و قابلیت اطمینان بستگی دارد. When data no longer fits on one machine, you can no longer do everything on one disk and one CPU. You need to distribute the data across multiple machines. There are three main approaches for this, compared in the table below. The choice depends on data volume, budget, and your need for speed and reliability.
| رویکردApproach | داده کجا ذخیره میشود؟Where is data stored? | نکات کلیدیKey points |
|---|---|---|
| خطلولهٔ محلی (یونیکس)Local pipeline (Unix) | روی دیسک همان ماشینOn the same machine's disk | مناسب برای: دادههای کوچک تا چند گیگابایت محدودیت: با بزرگتر شدن داده، کارایی افت میکندBest for: Small data up to a few gigabytes Limitation: Performance drops as data grows |
| سیستم فایل توزیعشده (HDFS)Distributed filesystem (HDFS) | فایل بین چند سرور تقسیم و کپی میشودFiles are split and replicated across servers | مزیت: اگر یک سرور خراب شود، داده از جای دیگر قابل بازیابی است مزیت: پردازش روی همان سروری انجام میشود که داده آنجاست (کاهش ترافیک شبکه)Advantage: If one server fails, data can be recovered from another Advantage: Processing happens on the same server where data lives (less network traffic) |
| ذخیرهسازی اشیاء (S3)Object storage (S3) | فایلها روی سرورهای راهدور و از طریق HTTP در دسترس هستندFiles are on remote servers, accessed via HTTP | مزیت: هزینهٔ ذخیرهسازی بسیار پایین و ماندگاری بالا محدودیت: تأخیر شبکه بیشتر از HDFS است، چون داده و پردازش از هم جدا هستندAdvantage: Very low storage cost and high durability Limitation: Higher network latency than HDFS because compute and storage are separate |
فرض کنید یک روز لاگ، ۸۰ گیگابایت حجم دارد. اگر بخواهید این فایل را با دستور sort روی لپتاپ خودتان مرتب کنید، ساعتها زمان میبرد. اما اگر یک خوشه (چندین ماشین که با هم کار میکنند) در اختیار داشته باشید، فایل را به قطعههای کوچکتر (مثلاً ۱۲۸ مگابایتی — همان اندازهٔ پیشفرض بلوک در HDFS) تقسیم میکنید و هر قطعه را روی یک ماشین جداگانه با grep پردازش میکنید. در نهایت، نتایج همهٔ ماشینها را با هم جمع میکنید. این دقیقاً همان ایدهٔ خطلولهٔ یونیکس است، فقط در مقیاس بسیار بزرگتر.
Imagine one day of logs is 80 GB. If you try to run sort on your laptop, it takes hours. But if you have a cluster (several machines working together), you split the file into smaller pieces (e.g. 128 MB each), process each piece on a separate machine with grep, and finally merge all the results. This is exactly the same Unix pipeline idea, just at a much larger scale.
خبر خوب این است که برای استفاده از این مدل، مجبور نیستید خودتان Hadoop یا سیستمهای مشابه را پیادهسازی کنید. در دنیای داتنت هم همین الگو را دارید: فایلهای روزانه در Blob Storage و یک BackgroundService یا Azure Function زمانبندیشده. برنامه شما فایل را تکهتکه میخواند (با PipeReader یا IAsyncEnumerable)، هر تکه را پردازش میکند و با SqlBulkCopy در مقصد مینویسد. فقط یادتان باشد: هرگز کل یک فایل بزرگ را یکجا در RAM بارگذاری نکنید.
The good news is that you don't need to implement Hadoop or similar systems yourself to use this model. In the .NET world, you already have the same pattern with daily files in Blob Storage and a scheduled BackgroundService or Azure Function. Your app reads the file in chunks (using PipeReader or IAsyncEnumerable), processes each chunk, and writes with SqlBulkCopy. Just remember: never load a large file entirely into RAM at once.
در بخش قبل دیدیم که وقتی داده از ظرفیت یک ماشین فراتر میرود، باید آن را میان چند ماشین پخش کنیم؛ اما روی چنین خوشهای چطور برنامه بنویسیم؟ پاسخ کلاسیک، MapReduce است: مدلی برنامهنویسی برای پردازش دادههای حجیم روی خوشهای از ماشینها. فرض کنید میلیونها سفارش فروش دارید و میخواهید بدانید هر کشور چقدر فروش داشته است. در این مدل، کار را در سه مرحله انجام میدهید: In the previous section we saw that when data outgrows one machine, you spread it across many. But how do you program such a cluster? The classic answer is MapReduce: a programming model for processing massive data on a cluster of machines. Imagine you have millions of sales orders and want to know each country's total sales. In this model, you do the work in three steps:
نگاشت (Transform): هر سفارش را به یک جفت (کلید, مقدار) تبدیل میکنید.
مثلاً سفارشی از ایران با مبلغ ۱۲۰ تومان میشود (IR, 120).
این مرحله کاملاً موازی است و هر ماشین بدون نیاز به هماهنگی با دیگران، بخشی از داده را پردازش میکند.
Map: Turn each order into a (key, value) pair.
For example, an order from Iran with amount 120 becomes (IR, 120).
This step is fully parallel — each machine processes its portion independently.
مرتبسازی و جابهجایی: حالا باید همهٔ جفتهایی که کلید یکسان دارند (مثلاً همهٔ IRها)
به یک ماشین منتقل شوند تا بتوان جمع هر کشور را محاسبه کرد.
این مرحله شامل انتقال داده روی شبکه است و معمولاً گرانترین و زمانبرترین بخش پردازش محسوب میشود.
Sort and move: Now all pairs with the same key (e.g., all IRs)
must be sent to one machine so we can calculate the total per country.
This involves moving data across the network and is often the most expensive and time-consuming part.
جمعآوری: در این مرحله، هر ماشین مقادیر مربوط به کلید خود را جمع میکند.
مثلاً همهٔ مبالغ ایران با هم جمع میشوند: IR → 120 + 90 + 40 + 80 = 330.
خروجی نهایی، فروش هر کشور است که معمولاً حجم بسیار کمتری نسبت به ورودی دارد.
Aggregate: In this final step, each machine sums the values for its key.
For example, all amounts for Iran are added up: IR → 120 + 90 + 40 + 80 = 330.
The final output is each country's total sales, which is usually much smaller than the input.
Extract (استخراج): داده از مبدأ خوانده میشود (فایل، دیتابیس، API).
Transform (تبدیل): داده پاکسازی، فیلتر یا تغییر شکل داده میشود.
Load (بارگذاری): دادههای تبدیلشده در مقصد نوشته میشوند (انبار داده، جدول تحلیلی).
مناسب برای: خطلولههای سنتی و انتقال داده بین سیستمها.
Extract: Data is read from a source (file, database, API).
Transform: Data is cleaned, filtered, or reshaped.
Load: The transformed data is written to a destination (data warehouse, analytics table).
Best for: Traditional pipelines and moving data between systems.
Map: هر رکورد به یک جفت کلید/مقدار تبدیل میشود.
Shuffle: همهٔ مقادیر یک کلید به یک ماشین منتقل میشوند.
Reduce: مقادیر هر گروه جمعآوری و محاسبه میشوند.
مناسب برای: پردازش موازی حجم عظیم داده روی خوشه.
Map: Each record is turned into a key/value pair.
Shuffle: All values for a key are sent to one machine.
Reduce: Values in each group are aggregated and computed.
Best for: Parallel processing of massive data on a cluster.
در دموی زیر، ۱۲ سفارش نمونه با دو روش مختلف پردازش میشوند.
شما میتوانید حالت پردازش را انتخاب کنید:
🔹 دستهای (Batch): رکوردها در گروههای ۴تایی با هم حرکت میکنند (۳ سفر شبکهای).
🔹 جریانی (Stream): هر رکورد به تنهایی سفر میکند (۱۲ سفر شبکهای).
همچنین میتوانید مدل خطلوله را عوض کنید:
🔹 ETL: مراحل Extract، Transform و Load را نشان میدهد.
🔹 MapReduce: مراحل Map، Shuffle و Reduce را نشان میدهد.
اعداد پایین صفحه، تعداد انتقالهای شبکه، تعداد رکوردهای پردازششده، بار هر سفر و زمان مفهومی را نشان میدهند تا بتوانید هزینهٔ هر روش را مقایسه کنید.
In the demo below, 12 sample orders are processed in two different ways.
You can choose the processing mode:
🔹 Batch: Records travel in groups of 4 (3 network trips).
🔹 Stream: Each record travels alone (12 network trips).
You can also switch the pipeline model:
🔹 ETL: Shows Extract, Transform, and Load stages.
🔹 MapReduce: Shows Map, Shuffle, and Reduce stages.
The numbers at the bottom show network trips, processed records, payload per trip, and conceptual time so you can compare the cost of each method.
۱۲ سفارش نمونه داریم. در حالت دستهای، رکوردها در گروههای ۴تایی با هم حرکت میکنند. در حالت جریانی، هر رکورد به تنهایی سفر میکند. داده و مسیر یکسان است، اما تعداد سفرهای شبکه در حالت جریانی ۴ برابر میشود. Twelve sample orders. In batch mode, records travel together in chunks of 4. In stream mode, each record travels alone. Same data, same path — but the number of network trips quadruples.
موتورهای جریان داده (Dataflow) مانند Spark و Flink داده را میان مرحلهها بیشتر در حافظه نگه میدارند و لازم نیست بعد از هر Map فایلهای بزرگی روی دیسک بنویسند. عملیاتهایی مثل Join و GroupBy همچنان به Shuffle نیاز دارند، اما API تمیزتری در اختیار برنامهنویس میگذارند؛ مثلاً بهصورت DataFrame یا زبان کوئری. Dataflow engines like Spark and Flink keep data in memory between stages more often, avoiding the need to write large files to disk after each Map. Operations like Join and GroupBy still use Shuffle, but with a cleaner API — for example, as a DataFrame or a query language.
# input: 12 paid orders
map(order) -> emit(order.Country, order.Amount)
# shuffle groups by key
IR: [120, 90, 40, 80]
DE: [80, 60, 70]
US: [200, 50, 150, 30, 110]
reduce(country, amounts) -> emit(country, sum(amounts))
# IR=330 DE=210 US=540
معادل سادهٔ این مدل در C#، کوئری LINQ زیر است:
orders.GroupBy(o => o.Country).Select(g => new { g.Key, Total = g.Sum(x => x.Amount) })
روی یک ماشین، این فقط یک کوئری حافظهای است. اما در یک خوشه، همان GroupBy تبدیل به یک Shuffle شبکهای میشود.
وقتی داده از حافظهٔ یک سرور فراتر برود، هزینهها کاملاً تغییر میکند.
The simple C# equivalent of this model is the LINQ query:
orders.GroupBy(o => o.Country).Select(g => new { g.Key, Total = g.Sum(x => x.Amount) })
On a single machine, it's just an in-memory query. But on a cluster, that same GroupBy becomes a network shuffle.
When data no longer fits in one server's RAM, the costs change entirely.
تا اینجا دیدیم پردازش دستهای چیست، چرا کار تحلیلی به ذخیرهسازی ستونمحور نیاز دارد و MapReduce چطور روی یک خوشه اجرا میشود. حالا نوبت به یک سؤال مهم میرسد: این دادهها دقیقاً کجا ذخیره میشوند؟ کار دستهای از یک جا میخواند و نتیجه را در یک جای دیگر مینویسد. فصل ۶ کتاب Fundamentals of Data Engineering سه مقصد اصلی را معرفی میکند. انتخاب هر کدام، تأثیر مستقیم بر هزینه، سرعت و نحوهٔ استفادهٔ بعدی از داده دارد. So far, we've seen what batch processing is, why analytical workloads need columnar storage, and how MapReduce runs on a cluster. Now it's time for an important question: where exactly is this data stored? A batch job reads from one place and writes results to another. Chapter 6 of Fundamentals of Data Engineering introduces three main destinations. Each choice directly impacts cost, speed, and how the data will be used later.
🔗 ارتباط با بخشهای قبلی: در بخش ۲ دیدیم که انبار داده معمولاً ستونمحور است. در بخش ۳ گفتیم که فایلها میتوانند روی HDFS یا S3 باشند. حالا در این بخش، این مفاهیم را کنار هم میچینیم و تفاوت این سه گزینه را با جزئیات بیشتری بررسی میکنیم. 🔗 Connection to previous sections: In section 2, we saw that data warehouses are usually column-oriented. In section 3, we mentioned that files can live on HDFS or S3. Now in this section, we bring these concepts together and examine the differences between these three options in more detail.
دادهای تمیز و منظم با ساختار مشخص. جدولها از قبل طراحی شدهاند، هر ستون نوع مشخصی دارد و همهچیز برای کوئریهای تحلیلی بهینه شده است. گزارشهای مالی، داشبوردهای مدیریتی و تحلیلهای کسبوکار معمولاً از اینجا تغذیه میشوند. انبار داده معمولاً بهصورت ستونمحور ذخیره میشود تا کوئریهای سنگین جمعآوری سریعتر اجرا شوند. Clean, structured data with a defined schema. Tables are designed in advance, each column has a specific type, and everything is optimized for analytical queries. Financial reports, management dashboards, and business analyses typically come from here. Data warehouses are usually column-oriented to speed up heavy aggregation queries.
دادهٔ خام، هر شکلی، هر ساختاری. فایلها به همان شکل اصلیشان در ذخیرهسازی ارزانقیمت (مثل S3 یا Azure Blob) قرار میگیرند. میتواند JSON، CSV، Parquet یا هر فرمت دیگری باشد. هزینهٔ ذخیرهسازی پایین است، اما تا وقتی داده را ساختاردهی و کاتالوگ نکنید، عملاً یک «باتلاق داده» خواهید داشت که پیدا کردن چیزی در آن سخت است. Raw data, any format, any structure. Files are stored in their original form in cheap storage (like S3 or Azure Blob). They can be JSON, CSV, Parquet, or any other format. Storage cost is low, but until you structure and catalog the data, it becomes a "data swamp" where finding anything is difficult.
ترکیبی از ارزانی دریاچه و ساختار انبار. داده همچنان بهصورت فایل در ذخیرهسازی ارزان نگهداری میشود، اما روی آن یک لایهٔ جدول، تراکنش و متادیتا اضافه میشود. هدف این است که هم هزینهٔ ذخیرهسازی پایین باشد و هم بتوان مثل یک انبار داده روی آن کوئری زد. این مدل در سالهای اخیر محبوبیت زیادی پیدا کرده است. A combination of lake cheapness and warehouse structure. Data is still stored as files in cheap storage, but a layer of tables, transactions, and metadata is added on top. The goal is to keep storage costs low while still being queryable like a data warehouse. This model has gained a lot of popularity in recent years.
فرض کنید هر شب، سفارشهای روزانه از SQL Server استخراج و بهصورت فایل Parquet در S3 (دریاچه) ذخیره میشوند. صبح روز بعد، یک کار دوم همان فایل را میخواند، پالایش میکند و در جدول fact_orders داخل انبار داده بارگذاری میکند. داشبوردهای مدیریتی فقط به انبار داده متصل هستند و گزارشهای خود را از آنجا میخوانند. فایل خام در دریاچه هم برای کارهایی مثل آموزش مدلهای یادگیری ماشین یا بازرسیهای بعدی نگهداری میشود. این ترکیب، هم انبار را تمیز نگه میدارد و هم دادهٔ خام را برای مواقع ضروری در دسترس قرار میدهد.
Imagine that every night, daily orders are extracted from SQL Server and stored as Parquet files in S3 (the lake). The next morning, a second job reads that file, cleans it, and loads it into the fact_orders table in the data warehouse. Management dashboards connect only to the warehouse and read their reports from there. The raw file in the lake is also kept for tasks like training machine learning models or future audits. This combination keeps the warehouse clean while keeping raw data available when needed.
در بسیاری از پروژهها، شما وظیفهٔ تولید فایل را دارید، نه مدیریت خود انبار داده. یک سرویسدهنده (مثل BackgroundService) داده را از منبع میخواند، فایل Parquet یا CSV را در یک مسیر مشخص در ذخیرهسازی اشیاء (مثل Azure Blob یا AWS S3) مینویسد و مسیر فایل را در یک جدول کنترل یا صف ثبت میکند تا کار بعدی آن را بردارد. یک نکتهٔ مهم: ذخیرهسازی اشیاء را مثل یک دیسک محلی فرض نکنید. فهرستکردن میلیونها فایل کوچک بسیار کند و گران است. بهتر است فایلهای روزانهای با حجم مناسب تولید کنید، نه تعداد زیادی فایل ریز.
In many projects, your job is to produce files, not to manage the data warehouse itself. A service (like a BackgroundService) reads data from the source, writes a Parquet or CSV file to a specific path in object storage (such as Azure Blob or AWS S3), and records the file path in a control table or queue for the next job to pick up. One important note: don't treat object storage like a local disk. Listing millions of tiny files is very slow and expensive. It's better to produce reasonably sized daily files, not a huge number of tiny ones.
در بخش قبلی دیدیم دادهٔ دستهای کجا ذخیره میشود؛ اما سؤال بعدی این است: خود داده چطور از سیستم مبدأ به آنجا میرسد؟ این فرایند را ورود داده یا Ingestion مینامند. ورود داده یعنی انتقال داده از جایی که تولید میشود (مثلاً دیتابیس عملیاتی) به جایی که کار دستهای میتواند آن را بخواند (مثلاً دریاچه یا انبار). این فرایند را نباید با یکپارچهسازی کامل سیستمها اشتباه گرفت؛ اینجا فقط داده جابهجا میشود و دو سیستم مستقل از هم میمانند. در این مرحله دو تصمیم کلیدی وجود دارد: In the previous section, we learned where batch data is stored. But another important question is: how does the data actually get from the source system to that destination? This process is called ingestion. Ingestion means moving data from where it's produced (e.g., an operational database) to where the batch job can read it (e.g., a lake or warehouse). This is different from full system integration; in ingestion, you're just transferring data, not connecting two systems. There are two key decisions here:
هر بار کل جدول را کپی میکنید. این روش بسیار ساده است و برای جدولهای کوچک (مثلاً چند هزار سطر) کاملاً مناسب است. اما برای جدولهای بزرگ با میلیونها سطر، هر شب کل تاریخ را دوباره خواندن، کاری اسرافآمیز و پرهزینه است. You copy the entire table every time. This is very simple and works well for small tables (e.g., a few thousand rows). But for large tables with millions of rows, re-reading all history every night is wasteful and expensive.
فقط سطرهای جدید یا تغییرکرده را میبرید. معمولاً از یک ستون زماندار مثل UpdatedAt > last_run یا قابلیت Change Tracking دیتابیس استفاده میکنید. این روش سبکتر است، اما حذفها و ناهماهنگی ساعت (Clock Skew) را باید خودتان مدیریت کنید.
You only take new or changed rows. Typically, you use a timestamp column like UpdatedAt > last_run or the database's Change Tracking feature. This is lighter, but you must handle deletes and clock skew issues yourself.
🔗 ارتباط با بخشهای قبلی: در بخش ۵ دیدیم که داده میتواند در انبار یا دریاچه ذخیره شود. حالا در این بخش میگوییم که داده چطور به آنجا میرسد. انتخاب بین Snapshot و Differential مستقیماً روی حجم دادهای که هر شب جابهجا میشود تأثیر میگذارد و در نتیجه، بر هزینه و زمان اجرای کارهای بعدی اثر دارد. 🔗 Connection to previous sections: In section 5, we saw that data can be stored in a warehouse or a lake. Now in this section, we explain how the data gets there. The choice between Snapshot and Differential directly affects the volume of data transferred each night, and consequently impacts cost and the execution time of subsequent jobs.
| ETL | ELT | |
|---|---|---|
| ترتیبOrder | استخراج → تبدیل → بارگذاریExtract → Transform → Load | استخراج → بارگذاری → تبدیلExtract → Load → Transform |
| تبدیل کجا انجام میشود؟Where is transform done? | قبل از بارگذاری در انبار، داخل کد شماBefore loading into the warehouse, in your code | داخل خود انبار، با قدرت SQL موتورInside the warehouse itself, using its SQL engine |
| چه موقع مناسب است؟When is it appropriate? | وقتی انبار قدرت پردازشی کافی ندارد، یا داده باید قبل از ورود تمیز و آماده شودWhen the warehouse isn't powerful enough, or data must be clean before loading | وقتی انبار قوی است و میخواهید دادهٔ خام را هم برای مصارف دیگر نگه داریدWhen the warehouse is powerful and you want to keep raw data for other purposes |
انتخاب اندازهٔ دستهها هم یک بدهبستان است: دستهٔ خیلی کوچک یعنی سربار (Overhead) زیاد برای هر بار commit به مقصد؛ دستهٔ خیلی بزرگ یعنی اگر کار نیمهکاره متوقف شود، باید حجم عظیمی را دوباره پردازش کنید. اندازهٔ مناسب به حجم داده و توان سیستم شما بستگی دارد. و برای تفاوت، همینقدر بدانید: جریان (Stream) بیپایان و بدون مرز است، اما دسته (Batch) همیشه مرز و حجم مشخصی دارد. Batch size is also a trade-off. A tiny batch means high commit overhead to the destination. A huge batch means that if the job dies midway, you'll have to reprocess a massive amount. Choosing the right size depends on your data volume and system capacity. For contrast only: a stream is unbounded, while a batch is always bounded.
var watermark = await control.GetLastSuccessAsync(jobName, ct);
var rows = await source.Orders
.Where(o => o.UpdatedAt > watermark)
.AsNoTracking()
.ToListAsync(ct); // استخراج تفاضلی
var clean = rows
.Where(o => o.Status == "Paid")
.Select(ToFactRow)
.ToList(); // تبدیل
await warehouse.BulkInsertAsync(clean, batchSize: 5000, ct);
await control.SaveWatermarkAsync(jobName, rows.Max(o => o.UpdatedAt), ct);
برای مدیریت این فرایند، یک جدول کنترل با ستونهای JobName، LastWatermark (نشانگر آخرین اجرای موفق) و Status طراحی کنید. کار شما باید ایدِمپوتنت (Idempotent) باشد، یعنی اجرای دوبارهٔ آن — حتی در همان روز — نباید سطر تکراری در مقصد بسازد. SqlBulkCopy با چند هزار سطر در هر دسته، معمولاً نقطهٔ شروع خوبی است. در مورد نحوهٔ انتقال: Push یعنی مبدأ داده را به سمت شما میفرستد؛ Pull یعنی کار شما خودش داده را از مبدأ میکشد. اکثر کارهای شبانه از روش Pull استفاده میکنند.
To manage this process, design a control table with JobName, LastWatermark (a marker of the last successful run), and Status. Your job should be idempotent, meaning that re-running it on the same day should not create duplicate rows in the destination. SqlBulkCopy with a few thousand rows per batch is usually a good starting point. Regarding the transfer method: Push means the source sends data to you; Pull means your job fetches the data from the source. Most nightly jobs use Pull.
تا اینجا با جزئیات فنی پردازش دستهای آشنا شدیم: از تعریف و ذخیرهسازی گرفته تا MapReduce و نحوهٔ ورود داده. حالا وقت آن است که ببینیم این مفاهیم در دنیای واقعی چگونه استفاده میشوند. فصل ۱۱ کتاب DDIA کاربردهای واقعی پردازش دستهای را مرور میکند و ما آنها را در چهار خانوادهٔ اصلی میبینیم. برای یک توسعهدهندهٔ بکاند، ETL آشناترین مورد است، اما بقیه نیز از همان الگوی کلی پیروی میکنند: ورودی بزرگ → خروجی مشتقشده. So far, we've learned the technical details of batch processing: from definition and storage to MapReduce and data ingestion. Now it's time to see how these concepts are used in the real world. Chapter 11 of DDIA reviews real-world batch use cases, which we group into four main families. For a backend developer, ETL is the most familiar, but the others follow the same general pattern: large input → derived output.
🔗 ارتباط با بخشهای قبلی: در بخش ۶ دیدیم که داده چطور از مبدأ به مقصد میرسد. حالا در این بخش میگوییم که پس از ورود داده، با آن چه کارهایی میتوان انجام داد. هر کدام از این کاربردها، یک «خروجی مشتقشده» تولید میکنند که معمولاً برای مصرف توسط سیستمهای دیگر (مانند داشبورد، API یا مدل یادگیری ماشین) آماده میشود. 🔗 Connection to previous sections: In section 6, we saw how data gets from source to destination. Now in this section, we explain what can be done with the data after it arrives. Each of these use cases produces a "derived output" that is typically ready for consumption by other systems (such as dashboards, APIs, or machine learning models).
استخراج، تبدیل، بارگذاری: سفارشها از دیتابیس عملیاتی (مثلاً SQL Server) استخراج میشوند، ارزها به یک واحد مشترک تبدیل میشوند، و در نهایت در جدول واقعیت (Fact Table) انبار داده بارگذاری میشوند. این همان خطلولهٔ شبانهای است که دادهٔ عملیاتی را به دادهٔ تحلیلی تبدیل میکند و برای گزارشگیری آماده میسازد. Extract, Transform, Load: Orders are pulled from the operational database (e.g., SQL Server), currencies are normalized to a common unit, and finally loaded into the fact table of the data warehouse. This is the nightly pipeline that turns operational data into analytical data, ready for reporting.
تبدیل داده به بینش: کوئریهای سنگین تحلیلی مثل جمع فروش هر کشور، محاسبهٔ نرخ بازگشت مشتریان، یا تحلیل قیف تبدیل (Conversion Funnel) روی انبار داده اجرا میشوند، نه روی همان دیتابیس عملیاتی. به این ترتیب، گزارشگیری سنگین، عملکرد سیستم عملیاتی را تحت تأثیر قرار نمیدهد. Turning data into insights: Heavy analytical queries such as sales by country, customer return rate calculation, or conversion funnel analysis run on the data warehouse, not on the operational database. This way, heavy reporting doesn't impact operational system performance.
ساخت ویژگیهای آموزشی: یک کار شبانه، ویژگیهای مورد نیاز مدل را از دادههای تاریخی استخراج میکند. مثلاً برای هر مشتری محاسبه میکند: «این مشتری در ۳۰ روز گذشته چند خرید انجام داده؟ مجموع مبلغ خریدش چقدر بوده؟» سپس این ویژگیها در یک فایل ذخیره میشوند و مدل روی آن فایل آموزش میبیند. این روش، دادههای خام را به دادههای آمادهٔ آموزش تبدیل میکند. Building training features: A nightly job extracts the required features for the model from historical data. For example, for each customer it calculates: "How many purchases has this customer made in the last 30 days? What's their total purchase amount?" Then these features are saved in a file and the model trains on that file. This method turns raw data into training-ready data.
آمادهسازی برای مصرف سریع: رتبهبندی پرفروشترین محصولات، کش پیشنهاد کالاهای مرتبط، یا هر نوع دادهٔ از پیش محاسبهشدهای که API باید سریع پاسخ دهد، یک بار در روز با یک کار دستهای ساخته میشود. سپس API فقط همان جدول یا کش آماده را میخواند و نیازی به محاسبهٔ سنگین در لحظه ندارد. Preparing for fast consumption: Bestseller rankings, recommendation caches, or any type of pre-computed data that the API needs to serve quickly is built once a day by a batch job. Then the API simply reads that ready table or cache, with no need for heavy real-time computation.
اگر مدیر کسبوکار بخواهد «فروش همین الان» را ببیند، یک کار شبانه پاسخگو نیست. در آن صورت باید سراغ پردازش جریانی (Streaming) بروید. اما واقعیت این است که بیشتر گزارشهای مدیریتی، تأخیر چند ساعته را به راحتی تحمل میکنند. پردازش دستهای در این موارد هم ارزانتر است و هم پیادهسازی سادهتری دارد. انتخاب بین دستهای و جریانی، همیشه به نیاز کسبوکار بستگی دارد. If a business manager wants to see "sales right now," a nightly job won't be enough. In that case, you need to look at streaming. But the reality is that most management reports can tolerate a few hours of delay. Batch processing is both cheaper and simpler to implement in these cases. The choice between batch and streaming always depends on the business requirement.
یک الگوی تمیز و قابلاعتماد در داتنت این است که API فقط وظیفهٔ ثبت سفارش را دارد و هیچ منطق گزارشگیری یا محاسبات سنگینی در آن اجرا نمیشود. یک BackgroundService یا IHostedService ساعت ۲ بامداد، جدول DailySales را از روی دادههای روز قبل محاسبه و بازسازی میکند. صفحهٔ مدیریت یا داشبورد، همین جدول آماده را میخواند. به این ترتیب، منطق گزارش در کنترلر پخش نمیشود، سرعت API حفظ میشود و کار دستهای مستقل و قابل اعتماد اجرا میشود.
A clean and reliable pattern in .NET is that the API only handles order registration and no reporting or heavy computation logic runs in it. A BackgroundService or IHostedService at 2 a.m. calculates and rebuilds the DailySales table from the previous day's data. The admin page or dashboard simply reads this ready table. This way, reporting logic doesn't leak into the controller, API performance is maintained, and the batch job runs independently and reliably.
اگر پاسخی را نمیدانید، به همان بخش برگردید. هدف، حفظکردن نام ابزارها نیست؛ باید بتوانید بگویید چرا یک کار دستهای به ذخیرهسازی و ریتم دیگری نیاز دارد.If you cannot answer one, return to that section. The goal is not tool names. You should be able to say why a batch job needs a different store and a different rhythm.