apache amoro is an open-source lakehouse management system that sits on top of table formats you already run — Apache Iceberg, Apache Paimon, and its own Mixed formats — and keeps them fast, tidy, and query-ready without you scheduling a single maintenance job. It began life as Arctic inside NetEase, was renamed Amoro when it entered the Apache Incubator, and its whole reason to exist is one stubborn problem: streaming and frequent-commit workloads bury a lakehouse table under thousands of tiny data files and delete files until reads crawl and metadata bloats. Amoro watches every table it manages and continuously compacts that mess in the background.
The key mental shift is that Amoro is not another table format competing with Iceberg — it is a management plane that operates Iceberg (and Paimon, and Mixed) tables for you. A long-running service called the Amoro Management Service (AMS) registers your catalogs, evaluates each table's file layout on a loop, and dispatches compaction work to a pool of worker processes. This guide walks the four ideas an interviewer will actually probe — the AMS and its catalog model, self-optimizing with its minor / major / full tiers, the optimizer and optimizer-group resource model, and the Mixed-Iceberg design for streaming upserts with fresh reads — and pairs each with a Solution-Tail interview answer: code, a step-by-step trace, an output table, then a concept-by-concept breakdown of why it works.
When you want hands-on reps immediately after reading, drill the optimization practice library →, rehearse partition-and-file-layout decisions on the partitioning practice set →, and harden your streaming-ingestion instincts on the streaming practice set →.
On this page
- Why Amoro changes lakehouse management in 2026
- The Amoro Management Service (AMS) & catalogs
- Self-optimizing — minor, major & full compaction
- The optimizer & optimizer-group model
- Mixed-Iceberg — streaming upserts & fresh reads
- Cheat sheet — Amoro recipes
- Frequently asked questions
- Practice on PipeCode
1. Why Amoro changes lakehouse management in 2026
Amoro is a management layer, not a table format — that one fact decides where it fits
The one-sentence invariant: Amoro does not store your data in a new format; it manages the Iceberg, Paimon, or Mixed tables you already have and keeps their file layout healthy automatically. Everything attractive about Amoro follows from that. There is no migration to a proprietary format, no new query engine to learn, and no lock-in at the storage layer — your tables remain readable by Spark, Flink, Trino, and any other engine that speaks the underlying format. Amoro adds an always-on service beside them.
The problem Amoro exists to solve — the small-files crisis.
- Streaming writers commit constantly. A Flink job writing every checkpoint (say every minute) produces a new data file — and for upsert tables, a new delete file — on every commit. After a day you have thousands of files, most a few kilobytes.
- Small files wreck reads. Every query must open, plan, and scan each file; task scheduling overhead dominates and throughput collapses. Metadata (manifests, snapshots) grows until even planning is slow.
- Delete files compound it. Merge-on-read formats keep equality- and position-delete files that the reader must apply on the fly. Left unmerged, a hot row can carry a long tail of deletes to reconcile at read time.
-
Manual maintenance is toil. On plain Iceberg you fight this by scheduling
rewrite_data_files,rewrite_manifests, andexpire_snapshotsyourself — as separate Spark jobs you size, monitor, and babysit per table.
What Amoro does instead.
- Self-optimizing. A background service watches each table's file statistics and continuously compacts fragments, merges delete files, and rewrites hot data — no cron, no hand-written maintenance jobs.
- One control plane. The Amoro Management Service registers catalogs, shows every table's optimizing status on a dashboard, and schedules the work across a shared pool of optimizers.
- Format-agnostic management. The same service manages native Iceberg, native Paimon, and Amoro's Mixed formats, so a heterogeneous lakehouse has one place to see and tune table health.
Where Amoro sits against the alternatives.
-
vs manual Iceberg maintenance.
rewrite_data_filesworks, but you own scheduling, sizing, back-pressure, and failure handling per table. Amoro turns that into a managed, evaluated-on-a-loop service with quotas and isolation. - vs catalog-only tools (Nessie, Polaris/Gravitino). Those govern metadata — branches, access, catalog federation. They do not compact your files. Amoro is complementary: it manages the physical health of tables.
- vs a warehouse (Snowflake/BigQuery). A managed warehouse hides file layout entirely, but you give up open storage and engine choice. Amoro keeps the open lakehouse and adds the automatic housekeeping the warehouse gave you for free.
What interviewers listen for.
- Do you say "Amoro manages Iceberg/Paimon tables, it is not a new format" in the first sentence? — senior signal.
- Do you name the small-files and delete-files problem as the reason it exists? — required framing.
- Do you contrast it with hand-scheduled
rewrite_data_filesrather than treating it as a query engine? — senior signal. - Do you mention self-optimizing and the AMS as the two load-bearing concepts? — the whole point.
Worked example — the file count that kills a streaming table
Detailed explanation. The clearest way to feel why Amoro exists is to count files. A minute-checkpoint Flink upsert job writes one data file and one delete file per commit into an Iceberg table. Nothing is wrong with any single commit — the damage is cumulative. Amoro's job is to keep the effective file count bounded no matter how often you write.
Question. A Flink job checkpoints every 60 seconds into a keyed Iceberg table, writing 1 data file + 1 equality-delete file per commit. With no compaction, how many files accumulate in 24 hours, and what does self-optimizing hold it near?
Input.
| parameter | value |
|---|---|
| checkpoint interval | 60 s |
| files per commit | 2 (1 data + 1 delete) |
| window | 24 h |
| target file size | 128 MB |
Code.
commits_per_day = (24 * 60 * 60) / 60 = 1440
files_per_day = commits_per_day * 2 = 2880 files/day (unmanaged)
with Amoro self-optimizing:
minor optimizing folds fragments into ~128MB segments continuously
effective_files ≈ table_bytes / 128MB + a small tail of recent fragments
Step-by-step explanation. Unmanaged, the table gains 2 files every minute — 2,880 per day, growing without bound — and each read must open all of them. With Amoro, a minor optimizing task fires whenever the count of fragment files crosses a trigger threshold, rewriting them into a handful of ~128 MB segment files and converting the equality-delete files into position deletes (and eventually merging them away). The steady-state file count tracks the table's real size, not its write frequency.
Output.
| after 24 h | unmanaged Iceberg | Amoro self-optimizing |
|---|---|---|
| data + delete files | ~2,880 and climbing | ~ table_size / 128 MB + small recent tail |
| read planning cost | O(files) — degrades daily | roughly constant |
| operator action | schedule + babysit rewrite jobs | none — AMS handles it |
Rule of thumb. If a table is written more often than once every few minutes, its long-term read cost is decided by who compacts it — schedule it yourself, or let Amoro's self-optimizing hold the file count near size / target.
2. The Amoro Management Service (AMS) & catalogs
AMS is the control plane — catalog management, the self-optimizing scheduler, a dashboard, and a terminal in one service
Amoro has one always-on process you must understand before anything else: the Amoro Management Service (AMS). An interviewer who asks "how is Amoro deployed?" wants the AMS and its responsibilities, then the catalog model that hangs off it. Get these crisp and the rest is configuration.
What the AMS is responsible for.
- Catalog management. AMS holds the registry of catalogs. Each catalog binds a metastore backend and a storage root, and every managed table lives under a catalog. You create catalogs in the AMS config or the dashboard.
- Self-optimizing scheduler. AMS periodically evaluates every managed table's file statistics, decides whether (and which kind of) optimizing is due, plans the tasks, and dispatches them to optimizers. This scheduling loop is the beating heart of Amoro.
- Dashboard & metrics. A web UI lists tables, shows optimizing status/history, file-size distributions, and resource usage — the single pane of glass for lakehouse health.
- SQL terminal. An in-dashboard terminal (backed by Spark/Kyuubi) to run ad-hoc SQL against managed tables without wiring up a separate client.
The catalog model.
-
Catalog type = format. A catalog is typed by the table format it manages:
iceberg,paimon,mixed-iceberg, ormixed-hive. The type tells Amoro how to read metadata and how to optimize. - Metastore backend. Under the catalog you choose where metadata lives: Hive Metastore, Hadoop (filesystem/HDFS catalog), AWS Glue, or a custom catalog implementation.
-
Storage & auth. The catalog also carries the warehouse root (
s3://,hdfs://, ...) and storage/auth config (S3 keys, Hadoop config, Kerberos) shared by all its tables.
The AMS system database — do not skip this in an interview.
- AMS stores its own state — catalog definitions, table runtime, optimizing task history, resource records — in a system database.
- The default embedded store is Derby, fine for a demo but single-node and non-durable.
- Production deployments point AMS at MySQL or PostgreSQL so the control plane survives restarts and can run HA. Forgetting to switch off Derby is the classic "it worked on my laptop" Amoro mistake.
Worked example — register an Iceberg catalog over Hive Metastore
Detailed explanation. Before Amoro can manage a table it must know the catalog. The everyday setup registers an Iceberg catalog whose metadata lives in Hive Metastore and whose data lives on S3. Once the catalog exists, every table under it becomes eligible for self-optimizing.
Question. Define an AMS catalog named prod_iceberg of type iceberg, backed by Hive Metastore at thrift://hms:9083, with a warehouse on s3://lake/warehouse. Show the catalog properties.
Input.
| field | value |
|---|---|
| catalog name | prod_iceberg |
| type | iceberg |
| metastore | hive |
| uri | thrift://hms:9083 |
| warehouse | s3://lake/warehouse |
Code.
name: prod_iceberg # amoro catalog definition (dashboard Catalogs -> Create)
type: iceberg # manages native Iceberg tables
table-formats: ICEBERG
storage-configs:
storage.type: S3
warehouse: s3://lake/warehouse
auth-configs:
auth.type: AK/SK
auth.access-key: ${S3_ACCESS_KEY}
auth.secret-key: ${S3_SECRET_KEY}
properties:
type: hive # metastore backend
uri: thrift://hms:9083
Step-by-step explanation. type: iceberg tells AMS to read Iceberg metadata and apply Iceberg-flavoured optimizing. The properties.type: hive plus uri point the catalog at Hive Metastore, so Amoro resolves tables through HMS exactly as Spark or Trino would. storage-configs give the warehouse root and S3 credentials shared by every table. Once saved, AMS lists all tables under prod_iceberg and begins evaluating them on its scheduling loop.
Output.
| AMS now knows | value |
|---|---|
| catalog |
prod_iceberg (type iceberg, metastore hive) |
| tables | every table registered in HMS under that warehouse |
| optimizing | each table evaluated on the scheduler loop once a group is set |
Rule of thumb. One catalog = one (format + metastore + storage) triple. Keep prod and dev in separate catalogs so their storage roots and credentials never bleed together.
Amoro interview question on catalog & metastore choice
Question. An interviewer says: "We already run Iceberg tables in AWS Glue and stream into them with Flink. We want Amoro to keep them compacted without moving the data. Which catalog type and metastore do you configure, and what stays unchanged for the readers?"
Solution Using a native-Iceberg catalog on Glue
Code.
name: analytics
type: iceberg # NATIVE iceberg — no format migration
table-formats: ICEBERG
storage-configs:
storage.type: S3
warehouse: s3://analytics/warehouse
properties:
type: glue # AWS Glue Data Catalog as the metastore
warehouse: s3://analytics/warehouse
# then, per table you want managed, run in the SQL terminal:
# ALTER TABLE analytics.db.events SET TBLPROPERTIES (
# 'self-optimizing.enabled' = 'true',
# 'self-optimizing.group' = 'default' );
Step-by-step trace.
| step | action | effect |
|---|---|---|
| 1 | register catalog type: iceberg, properties.type: glue
|
AMS reads the existing Glue-backed Iceberg tables in place |
| 2 | no data rewrite on registration | files, snapshots, schema untouched |
| 3 | set self-optimizing.enabled=true per table |
table enters the scheduler's evaluation loop |
| 4 | Flink keeps writing; Trino/Spark keep reading | readers still see plain Iceberg |
- Choosing
type: iceberg(notmixed-iceberg) means Amoro manages the existing Iceberg tables — there is no conversion, no new format. -
properties.type: gluepoints the catalog at the Glue Data Catalog, so table resolution matches what Flink and Trino already use. - Registration is metadata-only: Amoro does not touch data files until an optimizing task runs, and even then it only rewrites for compaction, preserving Iceberg semantics.
- Because the tables stay native Iceberg, every other engine reads them unchanged — Amoro is invisible to consumers except that reads get faster.
Output:
| concern | result |
|---|---|
| format | still native Iceberg (no migration) |
| readers (Trino/Spark) | unchanged, faster over time |
| what Amoro adds | continuous compaction + one dashboard |
Why this works — concept by concept:
-
Native-format catalog —
type: icebergmanages tables in place, so registering Amoro is a zero-risk, metadata-only step with no data movement. - Metastore parity — pointing the catalog at Glue (the same metastore the writers use) means Amoro, Flink, and Trino all resolve the identical table, avoiding split-brain metadata.
- Compaction, not conversion — optimizing only rewrites files for layout; it never changes the table's format contract, so downstream readers keep working.
-
Opt-in per table —
self-optimizing.enabledgates management table-by-table, so you roll Amoro out incrementally rather than all at once. - Cost — registration is O(1); ongoing management cost is the optimizer compute, bounded by the group's resources and each table's quota.
Optimization
Topic — optimization
Table maintenance and layout-optimization problems
3. Self-optimizing — minor, major & full compaction
Amoro compacts in three tiers — fragments to segments, delete-heavy segments, then the whole table
Self-optimizing is why Amoro exists, and interviewers drill it hard. The mechanism rests on one distinction — fragment files vs segment files — and three escalating kinds of compaction — minor, major, full. Say the file distinction first and the three tiers snap into place.
Fragment files vs segment files.
-
Target size.
self-optimizing.target-size(default 128 MB) is the size Amoro aims each output file toward. -
Fragment threshold. Files smaller than
target-size / self-optimizing.fragment-ratio(default ratio 8, so 16 MB) are fragment files — the small-file problem incarnate. - Segment files. Files at or above that threshold are segment files — already reasonably sized; they only need attention when they carry too many deletes.
The three optimizing types.
- Minor optimizing — cheap and frequent. Compacts many fragment files into segment files, and converts equality-delete files into position-delete files. This is the routine housekeeping that keeps a streaming table's file count flat.
- Major optimizing — rewrite delete-heavy segments. Takes segment files that have accumulated too many position deletes (redundant data) and rewrites them, applying the deletes so the reader no longer pays for merge-on-read.
- Full optimizing — the reset. Rewrites all files in the table (or partition) into the optimal layout and drops every delete file. Expensive; you run it rarely or on a slow cadence.
The triggers you tune.
-
self-optimizing.minor.trigger.file-count(default 12) — minor fires when the fragment-file count crosses this. -
self-optimizing.minor.trigger.interval(default 3600000 ms = 1 h) — also fire on a time interval so trickle tables still get tidied. -
self-optimizing.major.trigger.duplicate-ratio(default 0.1) — major fires when a segment's redundant (deleted) data ratio exceeds 10%. -
self-optimizing.full.trigger.interval(default -1, disabled) — set a period (e.g. daily) to schedule full rewrites.
Two things that bite people.
- Minor does not remove deletes; major does. If reads are slow because of delete accumulation, tune the major trigger, not the minor one.
- Full is not free. A full optimize rewrites the whole table — schedule it for low-traffic windows and give its group enough resources, or it will contend with ingestion.
Worked example — turning self-optimizing on and sizing fragments
Detailed explanation. The everyday task is enabling self-optimizing on a table and choosing a target size. Smaller targets mean more, smaller files (better write latency, worse read); larger targets mean fewer, bigger files (better read, heavier compaction). The defaults are sensible; you override when a table's read pattern demands it.
Question. Enable self-optimizing on db.events, target 256 MB files, and make minor optimizing fire once 8 fragment files pile up. Show the table properties.
Input.
| property | chosen value |
|---|---|
self-optimizing.enabled |
true |
self-optimizing.group |
default |
self-optimizing.target-size |
268435456 (256 MB) |
self-optimizing.minor.trigger.file-count |
8 |
Code.
ALTER TABLE db.events SET TBLPROPERTIES (
'self-optimizing.enabled' = 'true',
'self-optimizing.group' = 'default',
'self-optimizing.target-size' = '268435456', -- 256 MB
'self-optimizing.fragment-ratio' = '8', -- fragment < 32 MB
'self-optimizing.minor.trigger.file-count' = '8'
);
Step-by-step explanation. Setting enabled=true and a group puts the table into the AMS scheduler's evaluation loop. With target-size at 256 MB and fragment-ratio 8, any file below 32 MB counts as a fragment. Each scheduler pass counts fragment files; when the count reaches 8, AMS plans a minor optimizing task that rewrites those fragments into ~256 MB segments and rewrites equality deletes as position deletes. Segments stay untouched until a major trigger fires.
Output.
| behaviour after change | value |
|---|---|
| fragment threshold | 32 MB (= 256 MB / 8) |
| minor fires when | ≥ 8 fragment files present |
| output file target | ~256 MB segments |
| equality deletes | rewritten to position deletes on minor |
Rule of thumb. Tune target-size to your read pattern and minor.trigger.file-count to your write frequency — a fast writer wants a lower trigger so fragments never pile high.
Amoro interview question on the streaming small-files problem
Question. A Flink upsert job into an Iceberg table has made reads 5× slower over a week. File count is huge and there are thousands of equality-delete files. Which Amoro optimizing types fix each symptom, and what triggers do you set so it never regresses?
Solution Using minor + major triggers tuned to the write rate
Code.
ALTER TABLE db.orders SET TBLPROPERTIES (
'self-optimizing.enabled' = 'true',
'self-optimizing.group' = 'streaming',
'self-optimizing.target-size' = '134217728', -- 128 MB
-- keep fragment count low: fire minor early and often
'self-optimizing.minor.trigger.file-count' = '12',
'self-optimizing.minor.trigger.interval' = '600000', -- 10 min
-- attack delete accumulation: fire major when 10% is redundant
'self-optimizing.major.trigger.duplicate-ratio' = '0.1',
-- weekly reset for optimal layout
'self-optimizing.full.trigger.interval' = '86400000' -- daily
);
Step-by-step trace.
| symptom | cause | optimizing type | trigger set |
|---|---|---|---|
| thousands of tiny data files | per-checkpoint commits | minor | file-count 12, interval 10 min |
| thousands of equality-delete files | upsert deletes | minor (eq→pos) then major | duplicate-ratio 0.1 |
| segments carrying stale deleted rows | accumulated position deletes | major | duplicate-ratio 0.1 |
| long-term layout drift | continuous writes | full | daily interval |
- Minor collapses the fragment data files into 128 MB segments and turns equality-delete files into position deletes — this alone reverses most of the file-count explosion.
- Lowering
minor.trigger.intervalto 10 minutes means a fast writer never lets fragments pile beyond a few minutes' worth. - Major fires when a segment's redundant (deleted) fraction passes 10%, rewriting it to physically apply the deletes so the reader stops paying merge-on-read cost.
- Full on a daily interval periodically rewrites everything into the ideal layout and drops all delete files, catching anything the incremental tiers missed.
Output:
| metric | before | after tuning |
|---|---|---|
| effective file count | thousands, climbing | ~ size / 128 MB, flat |
| equality-delete files | thousands | folded into segments |
| read latency | 5× baseline | back to baseline |
Why this works — concept by concept:
-
Fragment vs segment — the whole engine keys off this threshold; only fragments trigger minor, so sizing
target-size/fragment-ratiocorrectly is what keeps compaction targeted. - Minor for files, major for deletes — matching the optimizing type to the symptom (file count vs delete ratio) is the senior distinction; using minor alone would never clear delete-heavy segments.
- Interval + count triggers — pairing a count trigger with a time trigger tidies both bursty and trickle tables, so no write pattern escapes compaction.
- Full as a safety net — a periodic full optimize bounds worst-case layout drift, guaranteeing the table can never slowly rot below a floor.
- Cost — minor is O(fragment bytes) and cheap; major is O(segment bytes with excess deletes); full is O(table bytes) and rare — total compute stays proportional to churn, not to table size.
Optimization
Topic — optimization
Compaction and small-files optimization problems
4. The optimizer & optimizer-group model
Optimizers are the compute that runs compaction; groups isolate resources and pin tables
The AMS plans optimizing tasks but does not execute them — that is the job of optimizers. An interviewer asking "how does Amoro scale compaction?" wants the optimizer, the optimizer-group, and how tables are pinned to a group. This is the resource-management half of self-optimizing.
The optimizer.
- What it is. An optimizer is a long-running execution process that pulls planned optimizing tasks from AMS and runs the actual file rewrites. It has a configured parallelism (threads / task slots).
- Where it runs — the container. Optimizers run in a container: LOCAL (a process on the AMS host, great for dev), FLINK (a Flink job on a YARN/Kubernetes cluster), or an external/custom container. Spark-based execution is likewise supported as an engine for the rewrite work.
- Scaling out. More throughput = more optimizers, or more parallelism per optimizer. Compaction capacity is just the sum of the optimizers' slots.
The optimizer group.
- A named resource pool. An optimizer group groups optimizers under one name with shared container settings (e.g. how much memory/parallelism a Flink optimizer gets).
-
Isolation. Groups isolate compute: a heavy
batchgroup running full optimizes cannot starve a latency-sensitivestreaminggroup. -
Pinning tables. A table is bound to a group with
self-optimizing.group. Every optimizing task for that table is dispatched only to optimizers in that group.
Fairness — quota.
-
self-optimizing.quotasets a table's share of its group's resources (roughly, how many optimizer-cores it may consume on average). - The scheduler uses quota to fair-share a group across many tables, so one churny table cannot monopolize the pool and starve the others.
Failure modes interviewers probe.
-
No running optimizer in the group = no optimizing. If a table's
self-optimizing.groupnames a group with zero live optimizers, AMS plans tasks that never run and the table quietly degrades. The dashboard's "pending" backlog is the tell. - Under-provisioned group. If churn outpaces the group's slots, the optimizing backlog grows; you add optimizers or raise parallelism rather than tightening triggers.
Worked example — a Flink optimizer group for streaming tables
Detailed explanation. The everyday setup is one dedicated group for streaming tables, backed by a Flink optimizer so compaction runs on your existing cluster and scales with it. You define the group's container, then launch an optimizer into it, then pin tables.
Question. Create an optimizer group streaming on a Flink container, launch one optimizer with 4 slots into it, and pin db.orders to it. Show the config.
Input.
| item | value |
|---|---|
| group name | streaming |
| container | flink |
| optimizer parallelism | 4 |
| table pinned | db.orders |
Code.
containers: # 1) define a Flink container + the group in AMS config
- name: flinkCon
container-impl: org.apache.amoro.server.manager.FlinkOptimizerContainer
properties:
flink-home: /opt/flink
optimizer_groups:
- name: streaming
container: flinkCon
properties:
taskmanager.memory: 2048
job-manager.memory: 1024
-- 2) scale one optimizer (4 parallelism) into the group (dashboard or API)
-- Optimizing -> Optimizer Groups -> streaming -> Scale-Out (parallelism = 4)
-- 3) pin the table to the group
ALTER TABLE db.orders SET TBLPROPERTIES (
'self-optimizing.enabled' = 'true',
'self-optimizing.group' = 'streaming',
'self-optimizing.quota' = '0.5'
);
Step-by-step explanation. The containers block teaches AMS how to launch a Flink job; the optimizer_groups block names streaming and gives its Flink optimizers 2 GB task-manager memory. Scaling out starts a Flink optimizer with 4 task slots that registers back to AMS. Pinning db.orders with self-optimizing.group = streaming routes its compaction only to that group, and quota = 0.5 caps its average share so siblings in the group still get serviced.
Output.
| after setup | value |
|---|---|
group streaming
|
1 optimizer, 4 slots, Flink container |
db.orders |
pinned to streaming, quota 0.5 |
| dispatch | orders' minor/major tasks run only on that optimizer |
Rule of thumb. One group per latency class — a streaming group sized for constant small compactions, a separate batch group for scheduled full optimizes — so heavy rewrites never stall your fresh-write tables.
Amoro interview question on resource isolation
Question. Two tables share one optimizer group. A nightly full optimize on the huge events table starves the small, latency-critical orders table, whose fragment backlog balloons every night. How do you fix it with Amoro's resource model, without slowing either table's compaction when run alone?
Solution Using separate groups plus quota
Code.
optimizer_groups:
- name: streaming # small, always-on, low-latency
container: flinkCon
properties: { taskmanager.memory: 2048 }
- name: batch # big, for scheduled full optimizes
container: flinkCon
properties: { taskmanager.memory: 8192 }
-- latency-critical table -> its own always-on group
ALTER TABLE db.orders SET TBLPROPERTIES (
'self-optimizing.group' = 'streaming',
'self-optimizing.quota' = '1.0'
);
-- heavy table -> the batch group, full optimize scheduled off-peak
ALTER TABLE db.events SET TBLPROPERTIES (
'self-optimizing.group' = 'batch',
'self-optimizing.full.trigger.interval' = '86400000' -- daily
);
Step-by-step trace.
| table | old (shared) | new (isolated) | result |
|---|---|---|---|
orders |
competes with events' full optimize | own streaming group |
minor runs promptly, backlog stays flat |
events |
starves orders nightly | own batch group, more memory |
full optimize runs without touching orders |
| groups | 1 pool, contention | 2 pools, isolated | no cross-interference |
- The root cause is shared compute: a
fulloptimize is O(table bytes) and monopolizes a shared group for hours. - Splitting into a
streamingand abatchgroup gives each workload its own optimizer slots, so events' heavy rewrite can never consume orders' capacity. - Pinning each table with
self-optimizing.grouproutes its tasks to the right pool;quotafair-shares within a pool if more tables join later. - Because each group is sized for its workload (2 GB vs 8 GB task managers), neither table is slower than before when run in isolation — you removed contention, not capacity.
Output:
| concern | outcome |
|---|---|
orders fragment backlog |
flat, even during events' full optimize |
events full optimize |
completes off-peak on dedicated compute |
| interference | none — physically separate optimizer pools |
Why this works — concept by concept:
- Optimizer group isolation — separate named pools give each workload dedicated slots, converting a contention problem into two independent, well-sized problems.
-
Group pinning —
self-optimizing.groupdeterministically routes a table's tasks, so the scheduler never sends a heavy full optimize into the latency-critical pool. - Quota fair-share — within a pool, quota bounds any one table's average consumption, protecting siblings from a single churny table.
- Right-sized containers — memory/parallelism per group is tuned to its workload, so isolation costs nothing in single-table throughput.
- Cost — total compute is the sum of the groups' slots; you trade a little idle headroom for predictable, contention-free compaction latency.
Optimization
Topic — optimization
Resource-isolation and scheduling-optimization problems
5. Mixed-Iceberg — streaming upserts & fresh reads
Mixed-Iceberg adds a ChangeStore over an Iceberg BaseStore — streaming upserts land fast, reads merge both
Amoro's own Mixed-Iceberg format is what you reach for when native Iceberg's commit cadence is too slow for the freshness you need. An interviewer asking "when would you use Mixed-Iceberg instead of native Iceberg?" wants the BaseStore / ChangeStore split, primary keys, and merge-on-read. It is built on Iceberg — every store is an Iceberg table underneath — so you keep compatibility while gaining streaming semantics.
The two stores (plus an optional third).
- BaseStore — the optimized baseline. An Iceberg table holding the compacted, deduplicated current data. This is where self-optimizing folds everything down to.
- ChangeStore — the streaming changelog. An Iceberg table that captures the CDC stream — inserts, update-before/update-after, deletes — written at high frequency by Flink. This is where fresh writes land first, cheaply.
- LogStore — optional real-time layer. A Kafka/Pulsar topic double-written alongside the file stores for second-level latency; without it, freshness is minute-level off the ChangeStore.
Primary keys and upserts.
-
Mixed-Iceberg tables are keyed. You declare a
PRIMARY KEY, and Amoro treats writes as upserts on that key rather than blind appends. - Change semantics. The ChangeStore preserves row kinds and ordering (via file sequence) so a delete-then-insert for the same key resolves to the latest state.
Reads — merge-on-read.
- MergeOnRead. A read of a Mixed-Iceberg table merges the BaseStore (compacted history) with the not-yet-merged ChangeStore rows, applying keys and deletes, to return the fresh, correct current state.
- Freshness vs cost. The more that sits unmerged in the ChangeStore, the more the reader pays at query time — which is exactly why self-optimizing continuously folds the ChangeStore into the BaseStore.
Self-optimizing is load-bearing here.
- On native Iceberg, self-optimizing just compacts files. On Mixed-Iceberg it also merges ChangeStore into BaseStore, which is what bounds merge-on-read cost and keeps fresh reads fast. Turn optimizing off on a Mixed table and reads degrade as the ChangeStore grows.
The trade-off, stated plainly.
- Native Iceberg — simplest, broadest engine support, freshness bounded by commit interval; ideal when minute-plus latency is fine.
- Mixed-Iceberg — fresher (minute or, with LogStore, second-level) streaming upserts, at the cost of running Amoro's optimizing pipeline and using Amoro-aware readers for the merged view.
Worked example — a keyed Mixed-Iceberg table for Flink upserts
Detailed explanation. The everyday Mixed-Iceberg table declares a primary key so Flink can upsert into it and readers get current-state semantics. The DDL looks like ordinary Flink SQL; the PRIMARY KEY ... NOT ENFORCED is what turns appends into upserts against the ChangeStore/BaseStore pair.
Question. Create a Mixed-Iceberg user_profiles table keyed on user_id in an Amoro catalog, then upsert a changing profile from Flink. Show what each store holds before a merge.
Input.
| event | user_id | tier | source |
|---|---|---|---|
| insert | 7 | free | CDC |
| update | 7 | pro | CDC |
Code.
-- Amoro Mixed-Iceberg catalog registered in Flink as `amoro`
CREATE TABLE amoro.db.user_profiles (
user_id BIGINT,
tier STRING,
updated_at TIMESTAMP(3),
PRIMARY KEY (user_id) NOT ENFORCED -- makes writes upserts
) WITH (
'self-optimizing.enabled' = 'true',
'self-optimizing.group' = 'streaming'
);
-- Flink streaming upsert (CDC): insert then update for the same key
INSERT INTO amoro.db.user_profiles
SELECT user_id, tier, updated_at FROM cdc_source; -- emits +I(7,free) then +U(7,pro)
Step-by-step explanation. The PRIMARY KEY (user_id) NOT ENFORCED declares this a keyed Mixed-Iceberg table, so Flink emits change events (+I, -U/+U, -D) that Amoro writes to the ChangeStore. Before any optimizing, the ChangeStore holds both the insert and the update for user_id = 7; the BaseStore is still empty. A merge-on-read query resolves the two change rows by key + sequence to the latest state (pro). Self-optimizing then merges the ChangeStore into the BaseStore, after which reads hit mostly compacted data.
Output.
| store | contents for user_id 7 (pre-merge) |
|---|---|
| ChangeStore |
+I(7, free), +U(7, pro) (both change rows) |
| BaseStore | empty |
| merge-on-read result |
7, pro (latest by key) |
Rule of thumb. Declare the primary key at create time — you cannot retrofit upsert semantics onto an append-only table, and without the key Mixed-Iceberg is just Iceberg with extra machinery.
Amoro interview question on streaming upserts with fresh reads
Question. You need near-real-time reads of a dimension that Flink updates thousands of times a minute by id, and analysts must always see the latest state within a minute. Explain how Mixed-Iceberg delivers this and what keeps read latency from degrading as writes pour in.
Solution Using a keyed Mixed-Iceberg table with self-optimizing
Code.
CREATE TABLE amoro.db.dim_customer (
id BIGINT,
email STRING,
status STRING,
updated_at TIMESTAMP(3),
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'self-optimizing.enabled' = 'true',
'self-optimizing.group' = 'streaming',
'self-optimizing.target-size' = '134217728',
'self-optimizing.minor.trigger.file-count' = '12',
'self-optimizing.minor.trigger.interval' = '300000' -- 5 min
);
-- Flink upserts by id land in the ChangeStore continuously;
-- minor optimizing folds them into the BaseStore every few minutes.
Step-by-step trace.
| moment | ChangeStore | BaseStore | read result |
|---|---|---|---|
| t0: 3k upserts/min flow in | fills with change rows | last compacted state | base ⊕ change (merge-on-read) |
| t1: minor optimizing fires | drained into base | absorbs latest per key | mostly base, tiny change tail |
| t2: analyst queries | small recent tail | current per-key state | fresh (< 1 min old) |
- Flink writes land in the ChangeStore as keyed change rows immediately — writes are cheap and never block on compaction, so ingestion keeps up with thousands of upserts a minute.
- Merge-on-read always combines BaseStore + unmerged ChangeStore by primary key, so a query at any instant sees the latest state — freshness does not wait for optimizing.
- Minor optimizing on a 5-minute interval continuously folds the ChangeStore into the BaseStore, so the unmerged tail the reader must reconcile stays small and read latency stays flat.
- If optimizing fell behind (too few optimizers), the ChangeStore tail would grow and reads would slow — which is why the table is pinned to a properly-sized
streaminggroup.
Output:
| requirement | how it is met |
|---|---|
| sub-minute freshness | merge-on-read over Base + Change (LogStore for second-level) |
| high write rate | cheap appends to ChangeStore, no blocking |
| stable read latency | minor optimizing keeps the unmerged tail small |
Why this works — concept by concept:
- ChangeStore write path — decoupling fast keyed writes from compaction lets ingestion absorb thousands of upserts a minute without back-pressure on the writer.
- Merge-on-read — reading Base ⊕ Change by primary key guarantees latest-state freshness independently of whether optimizing has caught up.
- Self-optimizing as the governor — minor optimizing bounds the unmerged ChangeStore tail, which is the single lever that keeps read latency flat under sustained write load.
- Primary-key upserts — keying the table is what makes "latest state per id" well-defined; without it there is no correct merge.
- Cost — writes are O(changes) and cheap; read cost is O(base scan + unmerged tail), and optimizing keeps that tail small, so steady-state reads are near pure-base cost.
Streaming
Topic — streaming
Streaming upsert and change-stream problems
Cheat sheet — Amoro recipes
Register an Iceberg catalog on AMS (Hive Metastore).
name: prod_iceberg
type: iceberg
storage-configs: { storage.type: S3, warehouse: s3://lake/warehouse }
properties: { type: hive, uri: thrift://hms:9083 }
Turn on self-optimizing and pick a group.
ALTER TABLE db.t SET TBLPROPERTIES (
'self-optimizing.enabled' = 'true',
'self-optimizing.group' = 'default'
);
Tune fragment size and minor trigger.
ALTER TABLE db.t SET TBLPROPERTIES (
'self-optimizing.target-size' = '134217728', -- 128 MB
'self-optimizing.fragment-ratio' = '8', -- fragment < 16 MB
'self-optimizing.minor.trigger.file-count' = '12'
);
Force periodic full optimize (drop all deletes).
ALTER TABLE db.t SET TBLPROPERTIES (
'self-optimizing.full.trigger.interval' = '86400000' -- daily
);
Launch a Flink optimizer into a group.
optimizer_groups:
- name: streaming
container: flinkCon
properties: { taskmanager.memory: 2048 }
# then dashboard: Optimizing -> streaming -> Scale-Out (parallelism = 4)
Create a keyed Mixed-Iceberg table (Flink SQL).
CREATE TABLE amoro.db.dim (
id BIGINT, val STRING, PRIMARY KEY (id) NOT ENFORCED
) WITH ('self-optimizing.enabled' = 'true', 'self-optimizing.group' = 'streaming');
Optimizing type picker.
| Symptom | Optimizing type |
|---|---|
| Many small (fragment) data files | minor |
| Equality-delete files piling up |
minor (eq → pos) |
| Segments heavy with position deletes | major |
| Layout drift / want a clean reset | full |
Frequently asked questions
What is Apache Amoro (formerly Arctic)?
Apache Amoro is an open-source lakehouse management system, currently in the Apache Incubator, that sits on top of open table formats — Apache Iceberg, Apache Paimon, and Amoro's own Mixed formats. It was originally built at NetEase under the name Arctic and renamed Amoro on entering incubation. Its core job is to keep managed tables healthy — continuously compacting small files and delete files via self-optimizing — through an always-on service called the Amoro Management Service (AMS).
Does Amoro replace Iceberg or Paimon?
No. Amoro is a management layer, not a table format, so it manages your existing Iceberg and Paimon tables in place rather than replacing them. Registering a native-Iceberg catalog is a metadata-only step with no data migration, and every other engine (Spark, Flink, Trino) keeps reading the tables unchanged. Amoro only adds continuous compaction plus one dashboard and scheduler for table health.
What is self-optimizing in Amoro?
Self-optimizing is Amoro's background compaction service: the AMS scheduler evaluates each managed table's file layout on a loop and dispatches compaction tasks to optimizers, with no cron jobs or hand-written maintenance from you. It fixes the streaming small-files problem by folding tiny fragment files into properly sized segments and merging delete files away. You enable it per table with self-optimizing.enabled=true and pin it to an optimizer group.
What is the difference between minor, major, and full optimizing?
Minor optimizing is cheap and frequent: it compacts small fragment files into segment files and converts equality-delete files into position-delete files. Major optimizing rewrites segment files that have accumulated too many position deletes, physically applying the deletes so reads stop paying merge-on-read cost. Full optimizing rewrites the entire table (or partition) into the optimal layout and drops every delete file — it is expensive and run rarely, usually on a scheduled interval.
What is an optimizer group?
An optimizer group is a named pool of optimizers (the long-running processes that actually run compaction) with shared container and resource settings. Groups isolate compute so a heavy batch workload cannot starve a latency-sensitive streaming one, and you pin a table to a group with self-optimizing.group. You scale compaction throughput by adding optimizers or parallelism to a group, and self-optimizing.quota fair-shares a group across the tables assigned to it.
What is the Mixed-Iceberg format?
Mixed-Iceberg is Amoro's own format, built on Iceberg, that adds a ChangeStore over an Iceberg BaseStore to support high-frequency streaming upserts with fresh reads. Keyed writes from Flink land cheaply in the ChangeStore, self-optimizing continuously folds them into the compacted BaseStore, and merge-on-read combines both by primary key so queries always see the latest state. An optional LogStore (Kafka/Pulsar) adds second-level freshness on top of the minute-level file stores.
Practice on PipeCode
Pipecode.ai is Leetcode for Data Engineering — every Amoro idea above, from self-optimizing's minor/major/full compaction to the optimizer-group resource model and the Mixed-Iceberg ChangeStore + BaseStore merge, maps to a hands-on practice room where you reason about file layout, compaction, and streaming upserts against real graded inputs. PipeCode pairs each reading with 450+ DE-focused problems and a real-time scoring engine, so your answer to "how would you keep this streaming table's reads fast?" holds up under a senior interviewer's depth probes.
Practice optimization problems now →
Streaming-ingestion drills →





Top comments (0)