The first scaling problem in our student fitness platform appeared in Java memory.
Schools uploaded Excel files containing thousands of students’ test results. Each row passed through parsing, validation, scoring rules, and evaluation logic. During concentrated uploads, several large tasks could keep substantial amounts of intermediate data alive at the same time. Processing smaller batches through a queue brought the memory peak under control.
The next problem appeared after those batches finished.
Their results stayed in MySQL. Every new testing campaign added more records, and historical results remained available for students, teachers, and reporting. Over time, the workload grew from millions of rows into tens of millions.
The row count was not a magic threshold. The more important change was that the same data now served several increasingly different workloads: teachers browsing a testing campaign, students reviewing their history, and administrators comparing schools, grades, and years.
This article follows that project’s architectural questions and develops a concrete design around them. The focus is on access paths, consistency boundaries, and corrected results—not on an assumed limit for MySQL or an unmeasured performance gain from ClickHouse.
The access paths mattered more than the row count
The original data model had three central concepts: a student, a testing campaign, and the student’s result for a particular test item.
A single campaign could produce several results per student: height, weight, lung capacity, sprint time, standing long jump, flexibility, and endurance, among others. Historical campaigns multiplied that volume.
Initially, conventional SQL optimization was useful. We reviewed indexes, reduced unnecessary joins, narrowed selected columns, and avoided deep offset pagination where possible. These changes addressed inefficient queries without introducing another storage system.
But the growing workload exposed a more structural problem.
Teachers usually entered through a campaign. They needed to browse results for the students participating in that campaign:
SELECT data_id, user_id, project_id, project_result
FROM fitness_result
WHERE detect_id = :detect_id
AND data_id > :last_data_id
ORDER BY data_id
LIMIT :page_size;
Students entered through their own identity. They wanted results across several campaigns:
SELECT data_id, detect_id, project_id, project_result
FROM fitness_result
WHERE user_id = :user_id
ORDER BY detect_id DESC, data_id DESC
LIMIT :page_size;
These examples illustrate access paths; production pagination must follow the actual ordering and identifiers used by the application.
Within one database, separate indexes can support both patterns. Once data is distributed across instances, routing becomes another concern.
Sharding by detect_id makes campaign queries straightforward, but a student’s history may span many shards. Sharding by user_id makes student history straightforward, but a large campaign may require querying many user shards.
A campaign-based shard also introduces a possible hotspot: one unusually large or active campaign can concentrate traffic on one node. Choosing a shard key therefore requires examining both ordinary requests and the largest expected workload.
Sharding middleware can route queries and merge results. It cannot decide which access path deserves predictable latency or how much fan-out the product can tolerate.
Time partitioning addressed a different question. It could help prune time-bounded queries and manage historical data, but partitioning within one MySQL instance would not add independent CPU, memory, or storage bandwidth. It was useful for data organization, not a substitute for workload capacity planning.
The architectural decision was therefore about the requests we needed to serve, rather than a rule that “ten million rows requires sharding.”
Keep one authoritative record and make other views explicit
Supporting several access paths does not necessarily require every request to use the same physical representation.
One option is to store complete query-oriented copies: one organized by result identifier, another by campaign, and another by student. Another is to store only relationships—such as user_id → data_id—and fetch details from the authoritative records afterward.
The trade-off is practical. Complete copies consume more storage and require more update handling, but reduce lookup stages. Relationship tables are smaller, but retrieving details can still fan out across several shards.
For the concrete design in this article, the boundary is:
One authoritative MySQL record, with derived views for other access paths and for analytics.
That distinction resolves an ambiguity common in architecture diagrams. If the authoritative row, campaign copy, student copy, and outbox record reside on different MySQL instances, they cannot all be covered by an ordinary local transaction.
Instead, the authoritative shard commits its own business change and event together:

An event relay then updates the campaign view, student-history view, and ClickHouse asynchronously. Each destination maintains its own local consistency boundary.
The outbox belongs with the authoritative write, rather than in a separate central database that would create another dual-write problem. This is the central guarantee of the transactional outbox pattern: business data and the event describing that change are committed in the same local transaction.
The cost is that a derived view can briefly lag behind the authoritative record.
That delay needs product behavior. Immediately after a teacher corrects a result, the edit screen can show the committed response or read the authoritative record. A campaign report may display its refresh time and update shortly afterward.
Pretending every screen is immediately consistent would hide the trade-off rather than solve it.
Events also need stable identities and source-generated versions. A retry must carry the same business version, and an older event arriving late must not overwrite a newer result in a derived MySQL view.
If a result is removed, that removal needs an explicit event as well. A synchronization design that handles only inserts and updates will eventually retain records that no longer belong in reports.
Move analytical scans away from operational requests
Even well-routed operational queries do not make reporting inexpensive.
Administrators may ask for average results by test item, pass rates across schools, or trends spanning several academic years. These requests often read large portions of the dataset and combine dimensions that are not part of the operational shard key.
MySQL can execute aggregations. The concern is what those aggregations compete with: result imports, corrections, student lookups, and teacher-facing pages.
The purpose of introducing ClickHouse is to isolate that analytical workload and use a storage layout suited to it.
A report grouped by campaign and test item may need only a few columns. Column-oriented storage can read those columns without reading every field in the application’s row model. Repeated dimensions such as school, grade, and test item can also compress effectively.
However, a useful analytical table requires more than copying operational columns.
For historical reports, school_id and grade_id should have defined meanings. Do they describe the student today or the student at the time of the test? If the intended question is “How did this school perform that semester?”, replacing historical dimensions with a student’s current school would change the meaning of past reports.
The same applies to scores. A raw sprint time and a normalized score answer different questions. Comparing results across grades or genders may require the scoring-rule version used at the time. Combining different test items into one raw average would be meaningless because their units differ.
These choices are part of the data model. ClickHouse can aggregate quickly, but it cannot determine what a historical comparison is supposed to mean.
Corrected results need a version model
Fitness results are not always immutable. A student may retake a test, or a teacher may correct an import error.
Suppose a standing long jump result changes from 2.21 meters to 2.31 meters. MySQL can update the authoritative row. On the analytical side, an append-based version model makes the change explicit.
A simplified ClickHouse table could look like this:
CREATE TABLE fitness_result_versions
(
data_id UInt64,
detect_id UInt64,
detect_date Date,
user_id UInt64,
project_id UInt32,
school_id UInt64,
grade_id UInt32,
project_result Float64,
version UInt64,
is_deleted UInt8,
event_time DateTime64(3)
)
ENGINE = ReplacingMergeTree(version)
PARTITION BY toYYYYMM(detect_date)
ORDER BY (detect_id, user_id, project_id, data_id);
Here, data_id is globally unique. The fields in ORDER BY are treated as immutable identifiers, and detect_date is fixed for the record’s partition assignment.
That is an important contract. ReplacingMergeTree identifies duplicates using the entire sorting key, not just whichever column the application considers its primary key. Changing one of those fields produces a different key and requires explicit handling of the old record.
The two versions of our example would be:
| data_id | project_result | version | is_deleted |
|---|---|---|---|
| 900001 | 2.21 | 1 | 0 |
| 900001 | 2.31 | 2 | 0 |
The second row is another insert.
The source service advances the version as part of the authoritative update and includes it in the outbox event. The ClickHouse consumer does not calculate the next version from arrival order. Messages can arrive late, be replayed, or be delivered more than once.
The version contract should ensure that the same record and version always describe the same state. Two conflicting payloads with an identical version make “latest” ambiguous.
Partition placement must also remain stable. Partitioning by an update timestamp could place an old version in May and its replacement in June. Background merges do not merge those partitions together.
Query correctness cannot depend on merge timing
Background merging may remove older versions, but a plain query can still encounter several versions of one record. Correctness should not depend on a merge happening before someone opens a report.
For a bounded query, FINAL can apply the engine’s replacement logic at query time:
SELECT data_id, project_result
FROM fitness_result_versions FINAL
WHERE detect_id = 10001
AND user_id = 4
AND project_id = 3
AND is_deleted = 0;
With this table design, a deletion is represented by a higher-version row whose is_deleted value is one. The query must select the latest state before excluding deleted records.
Another approach is to select the latest version explicitly with argMax. For example, this query calculates school-level averages for each test item using the latest valid record:
SELECT
tupleElement(latest, 1) AS project_id,
tupleElement(latest, 2) AS school_id,
count() AS result_count,
avg(tupleElement(latest, 3)) AS avg_result
FROM
(
SELECT
data_id,
argMax(
tuple(
project_id,
school_id,
project_result,
is_deleted
),
version
) AS latest
FROM fitness_result_versions
WHERE detect_id = 10001
GROUP BY data_id
)
WHERE tupleElement(latest, 4) = 0
GROUP BY project_id, school_id;
The order of operations matters: first identify the latest state of each result, then calculate the report.
The inner filter uses an immutable campaign identifier. Filtering on a potentially corrected attribute, such as school, before selecting the latest version could accidentally retain an older matching version.
Neither FINAL nor argMax should be declared universally cheap or expensive. Their cost depends on the scanned range, sorting key, number of versions, and ClickHouse version. Representative queries should be measured with the actual data distribution.
Pre-aggregation must account for corrections
Precomputing reports is attractive when the same campaign statistics are requested repeatedly. But adding a materialized view to a versioned table does not automatically make its aggregates correct.
Consider a view that adds every inserted result to a running sum and count. When the corrected value of 2.31 arrives, it may be added alongside the original 2.21. The later replacement of the old row in the source table does not automatically subtract its earlier contribution from the aggregate.
ClickHouse incremental materialized views process inserted data blocks; they are not continuously recalculated views of the source table’s latest logical state.
For this workload, there are several reasonable approaches.
Reports with manageable query ranges can select the latest results before aggregating. Frequently accessed reports can be rebuilt for affected campaigns and published as complete snapshots. More demanding low-latency reporting may justify explicit adjustment events or version-aware aggregate states.
The choice depends on how frequently results change and how fresh reports must be.
A practical starting point for semester-based testing is to refresh affected campaign reports after batches of changes. Each completed report snapshot has a refresh timestamp and a defined input cutoff. Readers use the latest completed snapshot rather than a partially rebuilt result.
If new events arrive during the rebuild, the campaign remains eligible for another refresh. The system should not lose that signal when the first rebuild finishes.
Adjustment events can reduce recomputation, but they create another consistency problem: replaying the same subtraction and addition must not apply the correction twice. For relatively infrequent retests and corrections, rebuilding a bounded report may be easier to operate correctly.
Reconciliation closes the operational loop
An outbox makes failed delivery recoverable, but it does not eliminate incorrect payloads, consumer bugs, delayed processing, or mistakes during historical backfills.
The system needs a way to compare what should exist with what the analytical side currently contains.
A raw count() comparison is insufficient. MySQL may contain one current row while ClickHouse physically stores three versions of that row. Both can be behaving correctly.
Reconciliation should compare the same logical dataset: the latest non-deleted records, within the same campaign and synchronization cutoff.
Counts are a useful first check, but equal counts do not prove equal results. A more detailed comparison can examine business identifiers, source versions, and relevant values. A task-level maximum version is also insufficient: one record can be current while another is missing.
For example, if MySQL contains version 3 of a result and ClickHouse contains version 2, a repair process can republish the authoritative version. Because the repair preserves the source identifier and version, it uses the same processing rules as ordinary delivery.
Backfills should follow that rule too. Historical imports must not assign newer versions merely because they run later. Otherwise, an old snapshot could overwrite a recent correction.
Operational monitoring should make unfinished work visible: the age of the oldest outbox event, consumer lag, failed records, campaigns awaiting report refresh, and reconciliation differences.
These signals are more actionable than a general statement that “synchronization is healthy.”
The final architecture gives each component a specific responsibility. MySQL owns business state and local transactions. Query-oriented views serve predictable operational access paths. ClickHouse handles analytical scans and reports. Outbox events connect them, while versioning and reconciliation make recovery possible.
The connection to the earlier OOM problem is straightforward. Batching controlled how much work existed in memory at once. The database changes controlled where different kinds of work were performed and how their results were represented.
The useful question was never simply whether MySQL could hold ten million rows. It was whether one physical model should continue serving every access path, every correction, and every report—and what guarantees had to remain explicit once those responsibilities were separated.
