DZone
Thanks for visiting DZone today,
Edit Profile
  • Manage Email Subscriptions
  • How to Post to DZone
  • Article Submission Guidelines
Sign Out View Profile
  • Post an Article
  • Manage My Drafts
Newsletter
Log In / Join
Refcards Trend Reports
Events Video Library
Refcards
Trend Reports

Events

View Events Video Library

Zones

Culture and Methodologies Agile Career Development Methodologies Team Management
Data Engineering AI/ML Big Data Data Databases IoT
Software Design and Architecture Cloud Architecture Containers Integration Microservices Performance Security
Coding Frameworks Java JavaScript Languages Tools
Testing, Deployment, and Maintenance Deployment DevOps and CI/CD Maintenance Monitoring and Observability Testing, Tools, and Frameworks
Partner Zones Build AI Agents That Are Ready for Production
Culture and Methodologies
Agile Career Development Methodologies Team Management
Data Engineering
AI/ML Big Data Data Databases IoT
Software Design and Architecture
Cloud Architecture Containers Integration Microservices Performance Security
Coding
Frameworks Java JavaScript Languages Tools
Testing, Deployment, and Maintenance
Deployment DevOps and CI/CD Maintenance Monitoring and Observability Testing, Tools, and Frameworks
Partner Zones
Build AI Agents That Are Ready for Production

Just dropped: New 2026 “Cloud-Native Foundations” Trend Report. See how teams are tackling complexity, cost & reliability.

AI can investigate. Engineers still decide. See how both work together across incident response in this DZone + Datadog webinar on Oct. 29.

Related

  • Grounding AI Agents in Governed Data
  • Engineering Self-Healing SQL Pipelines With LLMs: Validation, Guardrails, and Safe Recovery
  • One Agent, Two Runtimes: Defining State Ownership Between Temporal and LangGraph
  • Pipelines on Fire: Why Your CI/CD Tools Are the New Cyber Battlefield

Trending

  • Meta Wants to Run Your Business With AI — Microsoft and Salesforce Have a New Rival
  • AI on Top of a Dysfunctional System
  • AWS 7R Migration Strategies: A Decision Framework for Engineering Teams
  • Your Cloud Diagram Is Already Out of Date: An Operating Model for Continuous Security Architecture
  1. DZone
  2. Data Engineering
  3. Databases
  4. Apache Phoenix: Global Secondary Indexes With Tunable Consistency

Apache Phoenix: Global Secondary Indexes With Tunable Consistency

Apache Phoenix has evolved to complement HBase with several features to enhance the capabilities of the distributed database.

By 
Viraj Jasani user avatar
Viraj Jasani
·
Oct. 06, 26 · Analysis
Likes (1)
Comment
Save
Tweet
Share
137 Views

Join the DZone community and get the full member experience.

Join For Free

Apache Phoenix provides an open-source SQL interface over Apache HBase, combining the power of NoSQL horizontal scaling and sharding with SQL simplicity for low-latency and high-throughput OLTP operations on petabyte-scale data. 

Phoenix complements HBase by providing capabilities such as Global Secondary Indexes (GSIs), Atomic and Conditional updates, change data capture (CDC) streams, Updatable views, and multi-tenancy support across tables, indexes, and views. Furthermore, it supports server-side push-down execution for complex OLAP joins and grouping operations.

Phoenix-supported GSIs are backed by separate HBase tables. For each of the N indexes on the given data table, there are a total of (N + 1) HBase tables actively serving reads and writes: N index tables and one data table. The consistency mode of the index determines when exactly the data written on the indexes are available for reads. Let’s first understand read consistency in a database.

For distributed databases, here is a high-level consistency model:

  1. Eventual consistency: A reader will see the correct data eventually, but not necessarily immediately after writes.
  2. Read-your-writes: You will see your own writes, but not necessarily other writes immediately.
  3. Causal consistency: If A caused B (e.g., B is a reply to A), everyone sees A before B.
  4. Linearizability: All clients see operations in the same total order, and that order is consistent with the order they actually completed in. Single-key reads where you want no stale data.
  5. Serializability: Concurrent transactions appear to have run in some sequential order. MVCC, Locks, Optimistic Concurrency Control: Multi-key transactions that must not see partial state.
  6. External consistency (strict serializability): serializability + linearizability combined: transactions are serialized in real-time order. Use synchronized clocks or expensive coordination.

For HBase and Phoenix, strong consistency refers to serializability. By default, Phoenix-supported GSIs are strongly consistent, i.e., as soon as the write to the data table completes, readers are guaranteed to read updated data from the corresponding GSIs. 

Phoenix implements strongly consistent GSIs using two-phase commit and read-repair. Each data table with zero or more indexes has an IndexRegionObserver coprocessor attached to all the regions of the table. As part of the region coprocessor hooks preBatchMutate() and postBatchMutateIndispensably(), mutations to the index tables are generated and executed as RPC calls.

Two-phase commit for strong consistency

Two-phase commit for strong consistency


Despite having strongly consistent GSIs, Phoenix now also implements eventually consistent GSIs to support a broader range of highly scalable, latency- and throughput-sensitive applications to provide predictable and consistent write latencies regardless of the number of GSIs created on the data table.

Some advantages of using the eventually consistent GSIs:

  1. High availability for the data table writes
  2. Highly predictable tail latencies for the data table writes
  3. Improved write throughput for both the data table and the indexes

Phoenix provides two approaches to implement the eventually consistent GSIs using the CDC indexes. The writes to the CDC index always remain strongly consistent.

Approach 1: CDC Index With Serialized Index Mutations

In this approach, IndexRegionObserver on the data table generates all index mutations, for both strongly consistent and eventually consistent GSIs alike, after reading the current state of the data table row. For each data table row update, it serializes the eventually consistent index mutations, combines them into a single proto document, and then writes the document as a new cell on the CDC index row. 

This step applies to both pre- and post-update hooks. On the pre-index updates phase, the proto document contains all eventually consistent index mutations as unverified row updates; whereas on the post-index updates phase, the new proto document contains all eventually consistent index mutations as verified row updates.

Synchronous CDC index update

Synchronous CDC index update



Each data table region runs a single-threaded CDC consumer for the purpose of executing the eventually consistent index mutation RPCs in the background. This lets each CDC consumer scan only the change records for that region or partition. The CDC consumer scans the CDC index to retrieve the index mutations as the proto document for the given row, deserializes and executes them as RPCs. To improve the write throughput of the GSI writes, it uses a configurable batch size to execute several eventually consistent GSI mutations using a single RPC network call.

CDC Index Cell Structure for Approach 1

CDC Index Cell Structure for Approach 1


This approach is optimized for read I/O. The CDC consumer does not have to scan data table rows corresponding to the CDC index rows. However, it requires additional write I/O on the CDC index.

Approach 2: CDC as Lightweight Uncovered Index

In this approach, the CDC index is used as merely an uncovered index, with a single column, the empty column value as an unverified byte. The row key of the index remains the same as in approach 1, i.e., PARTITION_ID() + PHOENIX_ROW_TIMESTAMP() + data table primary keys. Since the CDC index has only a single cell with a one-byte value, the write operation on the data table is quite lightweight in comparison to approach 1. Moreover, unlike approach 1, only the pre-index update phase is required to make updates on the CDC index.

Generating the mutations for the eventually consistent GSIs requires the pre-image and post-image of each update done on the data table. To generate CDC pre-image and post-image for the changes, this approach requires performing a raw scan on the data table within a specific time range. Therefore, this approach is optimized for write I/O at the expense of additional read I/O on the data table. This approach is enabled by default, with the value of the config “phoenix.index.cdc.mutation.serialize” as “false.” Configure it to “true” to enable approach 1.

CDC Consumer Lifecycle: Ancestor/Descendent Relationships

Startup Phase

When an HBase region opens, the IndexRegionObserver coprocessor creates an IndexCDCConsumer worker if the data table has eventually consistent indexes.

Complete Parent Regions

When a region splits, or multiple regions merge, the consumer associated with splitting or merging parent regions stops processing change logs. The consumer of the child regions must continue from where the parent consumers left off.

Process any ancestor regions that were not fully processed before processing the immediate parent regions using the depth-first-search algorithm.

Resume or Start the Current Region/Partition

Check if SYSTEM.IDX_CDC_TRACKER contains the last processed timestamp for the current region. If yes, this region was moved from one server to another. Resume processing change logs from that timestamp; start from the beginning.

Process CDC index records with configurable batch size:

  • Query pattern: SELECT /*+ CDC_INCLUDE(DATA_ROW_STATE) */ PHOENIX_ROW_TIMESTAMP(), "CDC JSON" FROM WHERE PARTITION_ID() = ? AND PHOENIX_ROW_TIMESTAMP() > ? AND PHOENIX_ROW_TIMESTAMP() < ? ORDER BY PARTITION_ID() ASC, PHOENIX_ROW_TIMESTAMP() ASC LIMIT ?

Batch updates on indexes:

  • Execute BatchMutation on each eventually consistent GSI, increasing the index write throughput regardless of the client’s original write size.
  • For instance, even if the client application updates 10 rows on the data table using a single RPC call, the default batch size of 500 would let the CDC consumer update 500 or fewer rows for the given eventually consistent GSI in a single RPC call. This increases the write throughput on the GSIs.

Let’s take an example to understand the ancestor-descendant relationship for the replay of change records:

  • Region A splits into regions B and C
  • Region B splits into regions D and E
  • Regions E and C merge into F
  • Region A is fully processed
  • Region B is fully processed
  • Region C is in progress
  • Region E is in progress
  • New/current live regions: D and F


Apache Phoenix Database

Opinions expressed by DZone contributors are their own.

Related

  • Grounding AI Agents in Governed Data
  • Engineering Self-Healing SQL Pipelines With LLMs: Validation, Guardrails, and Safe Recovery
  • One Agent, Two Runtimes: Defining State Ownership Between Temporal and LangGraph
  • Pipelines on Fire: Why Your CI/CD Tools Are the New Cyber Battlefield

Partner Resources

×

Comments

The likes didn't load as expected. Please refresh the page and try again.

  • RSS
  • X
  • Facebook

ABOUT US

  • About DZone
  • Support and feedback
  • Community research

ADVERTISE

  • Advertise with DZone

CONTRIBUTE ON DZONE

  • Article Submission Guidelines
  • Become a Contributor
  • Core Program
  • Visit the Writers' Zone

LEGAL

  • Terms of Service
  • Privacy Policy

CONTACT US

  • 3343 Perimeter Hill Drive
  • Suite 215
  • Nashville, TN 37211
  • [email protected]

Let's be friends:

  • RSS
  • X
  • Facebook