DEV Community

Cover image for Apache Amoro (was Arctic): Self-Optimizing Lakehouse Management for Iceberg & Paimon
Gowtham Potureddi
Gowtham Potureddi

Posted on

Apache Amoro (was Arctic): Self-Optimizing Lakehouse Management for Iceberg & Paimon

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.

PipeCode blog header for Apache Amoro (was Arctic) — bold white headline 'Apache Amoro' with subtitle 'Self-Optimizing Lakehouse · Iceberg · Paimon' and a stylised small-files-into-compacted-segments scene on a dark gradient with purple, green, orange, and blue accents and a small pipecode.ai attribution.

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


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, and expire_snapshots yourself — 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_files works, 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_files rather 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
Enter fullscreen mode Exit fullscreen mode

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, or mixed-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.

Iconographic Amoro AMS diagram — the Amoro Management Service card holding catalog management, a self-optimizing scheduler, a dashboard and a SQL terminal, wired to catalogs over Hive Metastore / Hadoop / Glue and down to Iceberg and Paimon tables.

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
Enter fullscreen mode Exit fullscreen mode

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' );
Enter fullscreen mode Exit fullscreen mode

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
  1. Choosing type: iceberg (not mixed-iceberg) means Amoro manages the existing Iceberg tables — there is no conversion, no new format.
  2. properties.type: glue points the catalog at the Glue Data Catalog, so table resolution matches what Flink and Trino already use.
  3. 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.
  4. 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: iceberg manages 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.enabled gates 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

Practice →

Reliability Topic — reliability Durable control-plane and metastore problems

Practice →


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.

Iconographic Amoro self-optimizing diagram — fragment files versus segment files on the left, and three lanes for minor (compact fragments, equality-delete to position-delete), major (rewrite segments heavy with position deletes), and full (rewrite everything, drop all deletes) optimizing.

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'
);
Enter fullscreen mode Exit fullscreen mode

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
);
Enter fullscreen mode Exit fullscreen mode

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
  1. 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.
  2. Lowering minor.trigger.interval to 10 minutes means a fast writer never lets fragments pile beyond a few minutes' worth.
  3. 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.
  4. 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-ratio correctly 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

Practice →

Partitioning Topic — partitioning Partition-layout and file-sizing problems

Practice →


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 batch group running full optimizes cannot starve a latency-sensitive streaming group.
  • 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.quota sets 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.group names 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.

Iconographic Amoro optimizer-group diagram — the AMS scheduler planning optimizing tasks, dispatching them to optimizers running on Flink, Spark and a local container, grouped into isolated optimizer groups that tables are pinned to with self-optimizing.group.

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
Enter fullscreen mode Exit fullscreen mode
-- 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'
);
Enter fullscreen mode Exit fullscreen mode

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 }
Enter fullscreen mode Exit fullscreen mode
-- 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
);
Enter fullscreen mode Exit fullscreen mode

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
  1. The root cause is shared compute: a full optimize is O(table bytes) and monopolizes a shared group for hours.
  2. Splitting into a streaming and a batch group gives each workload its own optimizer slots, so events' heavy rewrite can never consume orders' capacity.
  3. Pinning each table with self-optimizing.group routes its tasks to the right pool; quota fair-shares within a pool if more tables join later.
  4. 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.group deterministically 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

Practice →

Concurrency Topic — concurrency Concurrent-worker and resource-contention problems

Practice →


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.

Iconographic Amoro Mixed-Iceberg diagram — a Flink CDC stream writing upserts into a ChangeStore and an optional LogStore, self-optimizing merging the ChangeStore into an optimized BaseStore, and a merge-on-read reader combining both stores for fresh reads on a primary-key table.

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)
Enter fullscreen mode Exit fullscreen mode

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.
Enter fullscreen mode Exit fullscreen mode

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)
  1. 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.
  2. 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.
  3. 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.
  4. 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 streaming group.

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

Practice →

Real-time Topic — real-time-analytics Fresh-read and merge-on-read analytics problems

Practice →


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 }
Enter fullscreen mode Exit fullscreen mode

Turn on self-optimizing and pick a group.

ALTER TABLE db.t SET TBLPROPERTIES (
  'self-optimizing.enabled' = 'true',
  'self-optimizing.group'   = 'default'
);
Enter fullscreen mode Exit fullscreen mode

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'
);
Enter fullscreen mode Exit fullscreen mode

Force periodic full optimize (drop all deletes).

ALTER TABLE db.t SET TBLPROPERTIES (
  'self-optimizing.full.trigger.interval' = '86400000'       -- daily
);
Enter fullscreen mode Exit fullscreen mode

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)
Enter fullscreen mode Exit fullscreen mode

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');
Enter fullscreen mode Exit fullscreen mode

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)