Em cenários como finanças, logística e IoT, os sistemas geram grandes volumes de dados de séries temporais — registros de transações, dados de trajetória e logs de monitoramento. Analisar esses dados na escala de terabytes em tempo real é um desafio comum de desempenho. O PolarDB for PostgreSQL utiliza recursos como tabelas particionadas e armazenamento em camadas quente-frio para oferecer uma solução econômica no armazenamento massivo de dados de séries temporais. Além disso, o recurso de índice columnstore (IMCI) permite executar análises de alto desempenho em tempo real sobre esses dados sem pré-processamento complexo, liberando efetivamente o valor dos seus dados.
Visão geral da solução
Fluxo de trabalho
Gravação de dados: sua aplicação de negócios grava dados de séries temporais (por exemplo, registros de transações) em um cluster PolarDB for PostgreSQL.
Índice columnstore: crie um índice columnstore na tabela base. O PolarDB for PostgreSQL mantém automaticamente os dados colunares junto ao armazenamento por linhas. Em comparação com o armazenamento por linhas, o armazenamento colunar organiza os dados por coluna, proporcionando maior taxa de compressão e lendo apenas as colunas relevantes durante consultas agregadas, o que reduz a E/S.
Aceleração de consultas: o sistema roteia consultas analíticas (como agregações de candlestick) preferencialmente para o índice columnstore, seja pelo otimizador ou por meio de um
Hintexplícito. O mecanismo de consulta utiliza o armazenamento colunar e o processamento paralelo para varrer e agregar os dados, retornando o resultado em seguida.
Vantagens da solução
Simplicidade de uso: não é necessário refatorar a aplicação nem realizar processos complexos de ETL. Basta criar um índice columnstore na tabela base para acelerar as consultas analíticas de forma transparente.
Rico em recursos: oferece suporte nativo a tabelas particionadas e inclui um conjunto abrangente de funções para análise de séries temporais, como
time_bucket,firstelast, simplificando o desenvolvimento SQL.
Resultados
Volume de dados: 100 milhões de linhas abrangendo 2 dias (aproximadamente 50 milhões de linhas por dia).
Consulta de agregação de candlestick: 5 métricas dentro de uma janela de tempo especificada — preço máximo, preço mínimo, preço de abertura, preço de fechamento e volume total de transações.
Grau de paralelismo do índice columnstore: 8.
-
Tempo de consulta (em segundos):
Cenário
Agregação de candlestick por segundo
Agregação de candlestick por minuto
Agregação de candlestick por hora
Agregação de candlestick por dia
Agregação de dados completos (100 milhões de linhas)
3,41
0,95
0,93
0,91
Agregação de dados de 1 dia (~50 milhões de linhas)
1,88
0,82
0,81
0,76
Agregação de dados de 12 horas (~25 milhões de linhas)
0,89
0,55
0,53
N/A
Agregação de dados de 1 hora (~6 milhões de linhas)
0,41
0,39
0,37
N/A
Procedimento
Etapa 1: Preparar o ambiente
-
Verifique se a versão e a configuração do seu cluster atendem aos seguintes requisitos:
-
Versões do cluster:
PostgreSQL 14 (versão secundária do mecanismo 2.0.14.10.20.0 ou posterior)
PostgreSQL 15 (versão secundária do mecanismo 2.0.15.15.7.0 ou posterior)
PostgreSQL 16 (versão secundária do mecanismo 2.0.16.8.3.0 ou posterior)
PostgreSQL 17 (versão secundária do mecanismo 2.0.17.7.5.0 ou posterior)
NotaVocê pode visualizar a versão secundária do mecanismo no console ou executando a instrução
SHOW polardb_version;. Se a versão secundária do mecanismo não atender aos requisitos, atualize a versão secundária do mecanismo. -
Defina o parâmetro
wal_levelcomological. Essa configuração adiciona as informações necessárias para decodificação lógica ao write-ahead logging (WAL).NotaVocê pode definir o parâmetro wal_level no console. A modificação desse parâmetro reinicia o cluster. Planeje suas operações de negócios adequadamente e proceda com cautela.
A tabela de origem deve ter uma chave primária, e a coluna da chave primária deve ser incluída ao criar o índice columnstore. Recomenda-se usar o tipo de dados
SERIALouBIGSERIALpara a chave primária, pois isso melhora significativamente a eficiência da sincronização de dados.É possível criar apenas um índice columnstore por tabela.
-
-
Ative o recurso de índice columnstore.
O método para ativar o IMCI varia dependendo da versão secundária do mecanismo do seu cluster PolarDB for PostgreSQL:
Etapa 2: Preparar os dados
-
Esta solução utiliza uma tabela de registros de transações e simula a geração de 100 milhões de registros abrangendo aproximadamente 2 dias. O horário de negociação é das 08:00 às 16:00 todos os dias, produzindo cerca de 40 milhões de registros diários.
-- Transaction records table CREATE TABLE market_trades ( trade_id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY, -- Auto-increment primary key trade_ts TIMESTAMP, -- Transaction timestamp market_id VARCHAR, -- Market ID price DECIMAL, -- Transaction price amount DECIMAL, -- Transaction amount insert_ts TIMESTAMP -- System write timestamp ); INSERT INTO market_trades(trade_ts, market_id, price, amount, insert_ts) SELECT trade_ts, market_id, price, amount, trade_ts + (random() * 500)::INT * INTERVAL '1 millisecond' AS insert_ts FROM ( -- ======================== -- 1. Day 1 peak: 2025-06-01 8:00 - 16:00, 40 million rows -- ======================== SELECT '2025-06-01 08:00:00'::TIMESTAMP + (random() * 28800)::INT * INTERVAL '1 second' + -- 28800 seconds = 8 hours (random() * 1000)::INT * INTERVAL '1 millisecond' AS trade_ts, CASE WHEN random() < 0.6 THEN 'BTC-USDT' ELSE 'ETH-USDT' END AS market_id, CASE WHEN random() < 0.6 THEN 30000 + (random() * 1000) ELSE 2000 + (random() * 100) END AS price, random() * 10 + 0.1 AS amount FROM generate_series(1, 40000000) UNION ALL -- ======================== -- 2. Day 1 off-peak: 2025-06-01 16:00 - 2025-06-02 08:00, 10 million rows -- ======================== SELECT CASE WHEN random() < 0.5 THEN -- 16:00 - 24:00 '2025-06-01 16:00:00'::TIMESTAMP + (random() * 28800)::INT * INTERVAL '1 second' ELSE -- 00:00 - 08:00 (early morning of day 2) '2025-06-02 00:00:00'::TIMESTAMP + (random() * 28800)::INT * INTERVAL '1 second' END + (random() * 1000)::INT * INTERVAL '1 millisecond' AS trade_ts, CASE WHEN random() < 0.6 THEN 'BTC-USDT' ELSE 'ETH-USDT' END AS market_id, CASE WHEN random() < 0.6 THEN 30000 + (random() * 1000) ELSE 2000 + (random() * 100) END AS price, random() * 10 + 0.1 AS amount FROM generate_series(1, 10000000) UNION ALL -- ======================== -- 3. Day 2 peak: 2025-06-02 8:00 - 16:00, 40 million rows -- ======================== SELECT '2025-06-02 08:00:00'::TIMESTAMP + (random() * 28800)::INT * INTERVAL '1 second' + (random() * 1000)::INT * INTERVAL '1 millisecond' AS trade_ts, CASE WHEN random() < 0.6 THEN 'BTC-USDT' ELSE 'ETH-USDT' END AS market_id, CASE WHEN random() < 0.6 THEN 30000 + (random() * 1000) ELSE 2000 + (random() * 100) END AS price, random() * 10 + 0.1 AS amount FROM generate_series(1, 40000000) UNION ALL -- ======================== -- 4. Day 2 off-peak: 2025-06-02 16:00 - 2025-06-03 08:00, 10 million rows -- ======================== SELECT CASE WHEN random() < 0.5 THEN -- 16:00 - 24:00 '2025-06-02 16:00:00'::TIMESTAMP + (random() * 28800)::INT * INTERVAL '1 second' ELSE -- 00:00 - 08:00 (early morning of day 3) '2025-06-03 00:00:00'::TIMESTAMP + (random() * 28800)::INT * INTERVAL '1 second' END + (random() * 1000)::INT * INTERVAL '1 millisecond' AS trade_ts, CASE WHEN random() < 0.6 THEN 'BTC-USDT' ELSE 'ETH-USDT' END AS market_id, CASE WHEN random() < 0.6 THEN 30000 + (random() * 1000) ELSE 2000 + (random() * 100) END AS price, random() * 10 + 0.1 AS amount FROM generate_series(1, 10000000) ) AS data; -
Crie um índice columnstore na tabela de registros de transações.
CREATE INDEX idx_csi_market_trades ON market_trades USING CSI;
Etapa 3: Executar consultas de agregação de candlestick
Caso de uso: calcular candlesticks em janelas de tempo fixas.
Exemplo de aplicação: calcular preço máximo, preço mínimo, preço de abertura, preço de fechamento e volume total de transações por segundo.
Os exemplos a seguir calculam dados de candlestick por segundo, por minuto, por hora e por dia, respectivamente.
Second-level candlestick aggregation
-- Second-level candlestick aggregation
/*+ SET (polar_csi.enable_query on) */
SELECT
time_bucket('1 second', trade_ts) AS candle_ts, -- Data within 1 second
market_id,
MIN(price) AS low, -- Lowest price within 1 second
MAX(price) AS high, -- Highest price within 1 second
FIRST(price ORDER BY trade_ts) AS open, -- Opening price within 1 second
LAST(price ORDER BY trade_ts) AS close, -- Closing price within 1 second
SUM(amount) AS vol -- Total transaction volume within 1 second
FROM market_trades
WHERE trade_ts >= '2025-06-01 00:00:00' AND trade_ts <= '2025-06-02 00:00:00'
GROUP BY candle_ts, market_id
ORDER BY candle_ts, market_id;
Minute-level candlestick aggregation
-- Minute-level candlestick aggregation
/*+ SET (polar_csi.enable_query on) */
SELECT
time_bucket('1 minute', trade_ts) AS candle_ts, -- Data within 1 minute
market_id,
MIN(price) AS low, -- Lowest price within 1 minute
MAX(price) AS high, -- Highest price within 1 minute
FIRST(price ORDER BY trade_ts) AS open, -- Opening price within 1 minute
LAST(price ORDER BY trade_ts) AS close, -- Closing price within 1 minute
SUM(amount) AS vol -- Total transaction volume within 1 minute
FROM market_trades
WHERE trade_ts >= '2025-06-01 00:00:00' AND trade_ts <= '2025-06-02 00:00:00'
GROUP BY candle_ts, market_id
ORDER BY candle_ts, market_id;
Hour-level candlestick aggregation
-- Hour-level candlestick aggregation
/*+ SET (polar_csi.enable_query on) */
SELECT
time_bucket('1 hour', trade_ts) AS candle_ts, -- Data within 1 hour
market_id,
MIN(price) AS low, -- Lowest price within 1 hour
MAX(price) AS high, -- Highest price within 1 hour
FIRST(price ORDER BY trade_ts) AS open, -- Opening price within 1 hour
LAST(price ORDER BY trade_ts) AS close, -- Closing price within 1 hour
SUM(amount) AS vol -- Total transaction volume within 1 hour
FROM market_trades
WHERE trade_ts >= '2025-06-01 00:00:00' AND trade_ts <= '2025-06-02 00:00:00'
GROUP BY candle_ts, market_id
ORDER BY candle_ts, market_id;
Day-level candlestick aggregation
-- Day-level candlestick aggregation
/*+ SET (polar_csi.enable_query on) */
SELECT
time_bucket('1 day', trade_ts) AS candle_ts, -- Data within 1 day
market_id,
MIN(price) AS low, -- Lowest price within 1 day
MAX(price) AS high, -- Highest price within 1 day
FIRST(price ORDER BY trade_ts) AS open, -- Opening price within 1 day
LAST(price ORDER BY trade_ts) AS close, -- Closing price within 1 day
SUM(amount) AS vol -- Total transaction volume within 1 day
FROM market_trades
WHERE trade_ts >= '2025-06-01 00:00:00' AND trade_ts <= '2025-06-02 00:00:00'
GROUP BY candle_ts, market_id
ORDER BY candle_ts, market_id;
Observações sobre SQL
/*+ SET (polar_csi.enable_query on) */: força a consulta a usar o plano de execução do índice columnstore. Em alguns cenários, o otimizador pode estimar incorretamente que o armazenamento por linhas é mais eficiente. Utilize esteHintpara garantir que a consulta utilize o caminho columnstore.time_bucket(bucket_width, ts): função fornecida pelo Time Series Database que agrupa o timestamptspelo intervalo especificado embucket_width.