Table of contents
Open Table of contents
Background
The order-query module mainly queries two tables:
- The core order table, order, stores order ID, creation time, and contact_id, the contact’s primary key.
- The contact table stores phone numbers, names, country, state, city, email, postal code, street address, and other contact information.
Queries mainly combine partial matching on contact fields with an order-creation-time filter. The SQL is roughly:
select c.*, o.order_id
from order o
join contact c on c.contact_id = o.contact_id
where c.first_name like '%test%' and OTHER_FILTER_COMBINATIONS
order and contact currently contain billions of rows. We have already optimized their indexes extensively, leaving no further room for index-based improvements.
We therefore decided to use Elasticsearch to accelerate partial-match searches.
Q1: Why retain so much data? Why not clean it up? A1: It represents more than a decade of data, from the system’s first day until now. Another team is cleaning up the whole database, but query optimization cannot wait for that work to finish. Our first challenge is migrating billions of rows without errors.
Technical Approach
Our system is gradually migrating from Oracle to PostgreSQL, so we want this optimization to build on PostgreSQL. The overall approach:
- Replicate Oracle data to PostgreSQL in real time.
- Monitor changes to the relevant PostgreSQL tables.
- Combine data from those tables and write it to Elasticsearch.
- Use Elasticsearch partial matching to find matching contact_ids.
- Query the database with those contact_ids to obtain order_ids.
- Perform the final order query and return results to the frontend.
Q1: Why query the database again? Can all information be stored in Elasticsearch? A1: Real queries are complex, involving many tables and complicated relationships. Combining everything into one wide table is difficult, especially with over a decade of accumulated data. The SQL above is highly simplified. Although the overall flow is complex, the bottleneck is partial matching in the contact table. Replacing that step and obtaining contact_ids makes the subsequent query very fast. The optimized SQL becomes:
select c.*, o.order_id
from order o
join contact c on c.contact_id = o.contact_id
where c.contact_id in ()
The contact_id set comes from Elasticsearch.
Q2: Why not ClickHouse, which uses fewer resources? A2: We are gradually moving monitoring logs from Elasticsearch to ClickHouse and plan to replace Elasticsearch eventually. However, the team has limited ClickHouse best-practice experience, and it has not yet been used in production. We are very familiar with Elasticsearch, so we chose it for now.
Cross-Database Migration
We use paid commercial software Qlik to replicate Oracle data to PostgreSQL in real time. We are considering replacing it with an open-source CDC solution, but the technical details of synchronizing this volume of data are still experimental for our team and not ready for production. For an open-source demo, see Synchronizing Oracle to PostgreSQL with Kafka Connect and Debezium CDC and Maintaining Cache Consistency.
Elasticsearch Index and Analyzer Design
PUT _index_template/order-contact-template
{
"priority": 1,
"template": {
"settings": {
"index": {
"number_of_shards": "3",
"number_of_replicas": "1"
},
"analysis": {
"analyzer": {
"order_contact_keyword_exact": {
"type": "custom",
"tokenizer": "keyword",
"filter": [
"lowercase"
]
}
}
}
},
"mappings": {
"properties": {
"@timestamp": {
"type": "date"
},
"city": {
"type": "text",
"analyzer": "order_contact_keyword_exact"
},
"company": {
"type": "text",
"analyzer": "order_contact_keyword_exact"
},
"country_iso3": {
"type": "keyword"
},
"email": {
"type": "text",
"analyzer": "order_contact_keyword_exact"
},
"first_name": {
"type": "text",
"analyzer": "order_contact_keyword_exact"
},
"last_name": {
"type": "text",
"analyzer": "order_contact_keyword_exact"
},
"order_id": {
"type": "keyword"
},
"organization_id_ref": {
"type": "keyword"
},
"phone": {
"type": "text",
"analyzer": "order_contact_keyword_exact"
},
"postal_code": {
"type": "text",
"analyzer": "order_contact_keyword_exact"
},
"state": {
"type": "keyword"
},
"street_address": {
"type": "text",
"analyzer": "order_contact_keyword_exact"
},
"title": {
"type": "text",
"analyzer": "order_contact_keyword_exact"
}
}
},
"aliases": {}
},
"index_patterns": [
"order-contact-*"
]
}
- This fast search supports only the last 5 years of orders. Split indexes by year, for example order-contact-2021 and order-contact-2022, to simplify query optimization and historical-data cleanup.
- To match database LIKE behavior exactly, do not split query fields into tokens. Use Keyword analysis. We normalize all fields to lowercase, requiring a custom analyzer.
- Replace LIKE with wildcard queries. See wildcard examples in Elasticsearch Queries.
- Verify the analyzer and queries against expectations.
- We tested the custom analyzer with Chinese, English, Japanese, Korean, German, and emoji. Results met our expectations and matched database LIKE behavior.
- Set Elasticsearch document_id to order_id rather than allowing it to generate IDs.
- @timestamp is the order table’s creation_date.
Synchronizing Data
We synchronize in Java, although CDC would be the best approach.
Standardizing on UTC
Elasticsearch stores timestamps in UTC. Convert database timestamps to UTC before writing, accounting for transitions between winter time (EDT) and summer time (EST) as described here. The corresponding field type is OffsetDateTime.
document.setCreationDate(localDateTime
.atZone(ZoneId.of("America/New_York")
.withZoneSameInstant(ZoneOffset.UTC)
.toOffsetDateTime());
Queries below also need UTC conversion.
Synchronizing Historical Data
Synchronize historical data one month at a time.
Comparing Elasticsearch and Database Counts
SELECT
COUNT(O.ORDER_ID)
FROM
ORDER O
JOIN CONTACT C ON O.CONTACT_ID = C.CONTACT_ID
WHERE
O.CONTACT_ID IS NOT NULL
AND O.CREATION_DATE IS NOT NULL
AND O.CREATION_DATE >= '2021-01-01'::date
AND O.CREATION_DATE < '2021-02-01'::date;
Specify the timezone in the Elasticsearch query.
POST order-contact-2021/_count
{
"query": {
"range": {
"@timestamp": {
"time_zone": "America/New_York",
"gte": "2021-01-01",
"le": "2021-02-01"
}
}
}
}
When the Counts Differ
During testing, a colleague found fewer Elasticsearch documents than database rows. My first suspicion was duplicate writes overwriting documents, but using order_id as document_id made that unlikely. Reviewing the code revealed LIMIT/OFFSET pagination ordered by creation time. Multiple orders can share a creation time, making this non-unique ordering unpredictable. The following comes from the PostgreSQL documentation:
When using LIMIT, it is important to use an ORDER BY clause that constrains the result rows into a unique order. Otherwise you will get an unpredictable subset of the query’s rows
Changing ORDER BY to order_id resolved it.
Deep Pagination Causes High Database CPU
When a month contained over a million rows, database CPU rose as high as 90%. LIMIT/OFFSET was probably responsible; our page size was 5,000.
We split monthly synchronization into daily queries. The optimized code:
LocalDate startDate = LocalDate.of(year, month, 1);
LocalDate endDate = startDate.with(TemporalAdjusters.lastDayOfMonth());
for (LocalDate date = startDate; !date.isAfter(endDate); date = date.plusDays(1)) {
while (true) {
queryDatabaseThenWriteIntoElasticsearch()
if (noMoreDatas) {
break;
}
}
}
Also replace OFFSET with cursor-style keyset (seek) pagination: pass the last order_id from one page into the next query. The optimized query:
" where o.order_id > " + seek +
" and o.contact_id is not null and o.creation_date is not null" +
" and o.creation_date >= '" + currentDay + " 00:00:00'" +
" and o.creation_date < '" + nextDayStr + " 00:00:00'" +
" order by o.order_id LIMIT " + PAGE_SIZE
Synchronizing Incremental Data
Capturing Changes
As noted earlier, Qlik replicates Oracle changes to PostgreSQL in real time. But Elasticsearch reads from PostgreSQL: how do we detect changes to order and contact? CDC is best. For the reasons above, we currently use database triggers and periodic Java polling.
Qlik’s replicated changes execute inserts/updates in PostgreSQL. Triggers capture those operations and write the changed data to a wide table.
create function insert_changed_contact()
returns trigger as
$$
declare
v_operation_type smallint;
v_contact_id bigint;
begin
if tg_op = 'INSERT' then
v_operation_type := 1;
v_contact_id := new.contact_id;
elsif tg_op = 'UPDATE' then
v_operation_type := 2;
v_contact_id := old.contact_id;
elsif tg_op = 'DELETE' then
v_operation_type := 3;
v_contact_id := old.contact_id;
end if;
insert into changed_contact (contact_id,
title,
city,
state,
zip,
company,
country_iso3,
organization_id_ref,
email,
first_name,
last_name,
address1,
day_phone,
operation)
values (v_contact_id,
new.title,
new.city,
new.state,
new.zip,
new.company,
new.country_iso3,
new.organization_id_ref,
new.email,
new.first_name,
new.last_name,
new.address1,
new.day_phone,
v_operation_type);
return new;
-- For BEFORE triggers, return NEW to allow the operation to proceed (potentially modified)
-- For AFTER triggers, return NEW (or OLD if appropriate)
end;
$$ language plpgsql;
create trigger contact_changed_trigger
after insert or update or delete
on contact
for each row
execute function insert_changed_contact();
A scheduled Java job polls new rows, writes them to Elasticsearch, then removes processed rows from the wide table.
Logstash Synchronization (Alternative)
This is essentially the Java approach implemented entirely in Logstash. We did not select it because we plan to move away from ELK and Logstash consumes substantial resources. Here is our implementation from that evaluation. These configurations were thoroughly verified; tune them for your situation.
Optimizing Logstash Configuration
- pipeline.id: eng-5361
# Logstash tries to load only files with .conf extension in the conf directory and ignores all other files, so we use .cfg extension
path.config: "/usr/share/logstash/pipeline/logstash-eng-5361.cfg"
# default value is 125
pipeline.batch.size: 100000
pipeline.workers: 1
pipeline
input {
jdbc {
jdbc_driver_library => "/usr/share/logstash/postgresql-42.5.0.jar"
jdbc_driver_class => "org.postgresql.Driver"
jdbc_connection_string => "${JDBC_CONNECTION_STRING}"
jdbc_user => "${JDBC_USER}"
jdbc_password => "${JDBC_PASSWORD}"
jdbc_default_timezone => "UTC"
schedule => "*/10 * * * * *"
statement_filepath => "/tmp/data-import/sql/delta.sql"
jdbc_fetch_size => 100000
jdbc_paging_enabled => true
jdbc_page_size => 100000
use_column_value => true
tracking_column_type => "numeric"
tracking_column => "unix_ts_in_secs"
last_run_metadata_path => "/tmp/data-import/last-value/order-contact-info-sql_last_value.yml"
tags => ["order_contact_info"]
}
}
filter {
if "order_contact_info" in [tags] {
mutate {
rename => {"order_creation_date" => "@timestamp"}
copy => { "order_id" => "[@metadata][order_id]" }
remove_field => ["order_id","contact_update_date","unix_ts_in_secs"]
}
}
}
output {
if "order_contact_info" in [tags] {
elasticsearch {
hosts => "elasticsearch:9200"
ssl => true
cacert => "config/elasticsearch-ca.pem"
user => "elastic"
password => "${ELASTIC_PASSWORD}"
index => "order-contact-%{+YYYY}"
action => "update"
doc_as_upsert => true
document_id => "%{[@metadata][order_id]}"
}
}
}
delta.sql Contents
SELECT c.*,extract(epoch from c.contact_update_date) AS unix_ts_in_secs
FROM contact c
WHERE extract(epoch from c.contact_update_date) > :sql_last_value
AND c.contact_update_date < LOCALTIMESTAMP
ORDER BY c.contact_update_date ASC, contact_id
Q1: Why convert time to a Unix timestamp and require it to be greater than the previous synchronized value but less than the current timestamp? A1: This official Elasticsearch blog post demonstrates periodically synchronizing MySQL data with Logstash and explains these conditions in detail.
order-contact-info-sql_last_value.yml stores the Unix timestamp of the final row from the previous synchronization. This configuration handles incremental data. Historical synchronization is similar and can be derived from the Logstash documentation.
Querying
Support queries against both Elasticsearch and the database so that searches remain available if Elasticsearch fails or is under maintenance.
Elasticsearch query optimizations:
- Query only the indexes required by the date range. For orders from 2022-01-01 to 2023-01-01, query order-contact-2022 and order-contact-2023.
- Use search_after for pagination, following the same idea as the database keyset approach.
- Split queries.
- If the index’s maximum max_date is later than the query end date, use Elasticsearch for the entire query.
- If it is earlier than the current date, split the query: [startDate, maxDateFromEs] goes to Elasticsearch; (maxDateFromEs, endDate] goes to the database.
- If max_date and endDate fall on the same day, still use the second approach, subtracting 1 day from max_date because we assume the synchronization gap does not exceed a day.
Production Elasticsearch Cluster Specifications
- Self-managed 3-node cluster.
- Each instance: Rocky Linux 8, ARM, 4 CPUs, 16 GB RAM, and 60 GB disk.
- We did not install Kibana. For convenient interaction, we use Elasticvue.