0Pricing
Pandas & NumPy Academy · Aula

Agregação incremental entre partes

Acumule contagens, somas e valores mínimos/máximos em execução entre as partes sem armazenar o arquivo inteiro na memória.

Agregação incremental entre partes é uma aula grátis de Pandas & NumPy Academy no CoddyKit. Esta é a aula 2 de 4. Você pode ler a aula completa abaixo gratuitamente — depois pratica ao vivo no navegador com um editor de código integrado e um tutor de IA 24/7. Faz parte do caminho de aprendizado de Pandas & NumPy Academy, e seu progresso é sincronizado entre a web e o app CoddyKit. O curso de Pandas & NumPy Academy inclui 4 aulas no total.

Por que usar agregação progressiva?

A agregação progressiva é fundamental para analisar conjuntos de dados maiores que a RAM sem distribuir a computação entre várias máquinas. Em vez de carregar todos os dados para calcular uma estatística final, você mantém acumuladores progressivos — somas parciais, contagens e valores mínimos/máximos — e os atualiza a cada bloco. O resultado final é montado a partir desses acumuladores leves depois que o arquivo inteiro é examinado. Esse padrão permite processar arquivos de terabytes em um único laptop.

Contando linhas e calculando a média

Calcular a média entre blocos exige acompanhar separadamente a soma e a contagem acumuladas. Você não pode simplesmente calcular a média das médias de cada bloco, pois os blocos podem ter tamanhos diferentes. A fórmula correta é total_sum / total_count. Esse padrão se estende a qualquer quantidade que possa ser decomposta: variância, correlação e histogramas também têm fórmulas progressivas.

import pandas as pd

total_sum = 0.0
total_count = 0

for chunk in pd.read_csv('transactions.csv', chunksize=100000):
    total_sum += chunk['amount'].sum()
    total_count += chunk['amount'].notna().sum()

grand_mean = total_sum / total_count
print(f'Rows processed: {total_count:,}')
print(f'Grand mean: {grand_mean:.4f}')

Mínimo e máximo progressivos

Acompanhar o mínimo e o máximo globais entre blocos é simples: inicialize com float('inf') e float('-inf') do Python e depois atualize com o mínimo/máximo de cada bloco. Isso evita qualquer armazenamento intermediário. O padrão pode ser generalizado para mínimo/máximo por grupo, mantendo um dicionário indexado pelo identificador do grupo.

import pandas as pd

global_min = float('inf')
global_max = float('-inf')

for chunk in pd.read_csv('prices.csv', chunksize=50000):
    chunk_min = chunk['price'].min()
    chunk_max = chunk['price'].max()
    if chunk_min < global_min:
        global_min = chunk_min
    if chunk_max > global_max:
        global_max = chunk_max

print(f'Price range: {global_min} to {global_max}')

Contagens progressivas de frequências

Para colunas categóricas, mantenha um dicionário progressivo de frequências, adicionando o resultado de value_counts() de cada bloco a um acumulador Series do Pandas. Como a adição de Series do Pandas alinha os rótulos dos índices, as categorias desconhecidas dos blocos posteriores são incluídas automaticamente. Depois de processar todos os blocos, ordene pela contagem para ver as principais categorias de todo o conjunto de dados.

import pandas as pd

freq = pd.Series(dtype='int64')

for chunk in pd.read_csv('orders.csv',
                         chunksize=100000,
                         usecols=['category']):
    chunk_counts = chunk['category'].value_counts()
    freq = freq.add(chunk_counts, fill_value=0)

# Final sorted frequency table
print(freq.sort_values(ascending=False).head(10))

Agregação GroupBy progressiva

Para calcular uma soma ou contagem com groupby entre blocos, aplique groupby().agg() dentro de cada bloco e armazene a Series ou o DataFrame resultante. Após o laço, concatene todos os resultados parciais e aplique um segundo groupby para combiná-los. Essa abordagem em duas etapas trata corretamente os grupos que aparecem em vários blocos, algo comum quando os dados são ordenados por data, e não por grupo.

import pandas as pd

partials = []
for chunk in pd.read_csv('sales.csv',
                         chunksize=100000,
                         usecols=['region', 'product', 'revenue']):
    p = chunk.groupby(['region', 'product'])['revenue'].sum()
    partials.append(p)

final = (
    pd.concat(partials)
    .groupby(level=['region', 'product'])
    .sum()
    .sort_values(ascending=False)
)
print(final.head(10))

Calculando a variância progressivamente (método de Welford)

Calcular a variância entre blocos é mais difícil que calcular a média. A fórmula ingênua E[X²] - E[X]² sofre de cancelamento catastrófico quando as médias são grandes. O algoritmo online de Welford mantém uma média progressiva e uma soma dos desvios quadráticos, atualizando-os a cada novo valor de maneira numericamente estável. Embora o SciPy implemente esse recurso, compreender o padrão permite estendê-lo à variância ponderada e à covariância.

import pandas as pd
import numpy as np

# Simple two-pass approach using stored chunk stats
chunk_stats = []
for chunk in pd.read_csv('data.csv',
                         chunksize=100000,
                         usecols=['value']):
    n = chunk['value'].count()
    mean = chunk['value'].mean()
    var = chunk['value'].var(ddof=1)
    chunk_stats.append((n, mean, var))

# Combine: use pooled variance formula
total_n = sum(s[0] for s in chunk_stats)
total_mean = sum(s[0]*s[1] for s in chunk_stats) / total_n
pooled_var = sum((s[0]-1)*s[2] + s[0]*(s[1]-total_mean)**2
                 for s in chunk_stats) / (total_n - 1)
print(f'Grand variance: {pooled_var:.4f}')

Criando um histograma progressivo

Calcular a distribuição de uma coluna em um arquivo grande demais para a RAM exige um histograma progressivo. Fixe previamente os limites das classes (com base em uma pequena amostra ou no conhecimento do domínio), depois use np.histogram(chunk_values, bins=edges) em cada bloco e acumule as contagens. No final, represente as contagens combinadas em um gráfico de barras. É assim que sistemas de processamento contínuo, como Kafka Streams e Flink, calculam histogramas aproximados.

import pandas as pd
import numpy as np

# Decide bin edges from a sample
sample = pd.read_csv('amounts.csv', nrows=5000)
bins = np.linspace(sample['amount'].min(),
                   sample['amount'].max(), 21)
counts = np.zeros(len(bins) - 1, dtype='int64')

for chunk in pd.read_csv('amounts.csv',
                         chunksize=100000,
                         usecols=['amount']):
    chunk_counts, _ = np.histogram(
        chunk['amount'].dropna(), bins=bins
    )
    counts += chunk_counts

print('Histogram counts:', counts[:5], '...')

Acompanhando valores únicos de forma aproximada

Contar valores distintos exatos entre blocos exige armazenar todos os valores únicos — potencialmente milhões. Para contagens aproximadas em grande escala, use um esboço HyperLogLog, disponível em Python por meio da biblioteca hyperloglog. Como alternativa, acompanhe os valores únicos de cada bloco com um conjunto e obtenha a união, mas isso cresce sem limite. Para uma aproximação econômica, use pd.Series.nunique() por bloco e informe a média — não é exato, mas geralmente é suficiente para criar o perfil dos dados.

import pandas as pd

unique_ids = set()
for chunk in pd.read_csv('events.csv',
                         chunksize=100000,
                         usecols=['user_id']):
    unique_ids.update(chunk['user_id'].dropna().unique())

print(f'Distinct user IDs: {len(unique_ids):,}')
# Warning: the set may grow large for high-cardinality columns

Relatando o progresso durante execuções longas

Processar um arquivo de vários gigabytes pode levar minutos. Adicione um relatório de progresso para saber se o fluxo de processamento está em execução e estimar o tempo restante. Conte os bytes ou as linhas processados e compare-os com o tamanho do arquivo. A biblioteca tqdm torna isso trivial com seu invólucro tqdm(reader). Mesmo sem tqdm, imprimir uma linha de status a cada 10 blocos fornece informações valiosas durante tarefas longas em lote.

import pandas as pd
import time

chunksize = 100000
start = time.time()
rows_processed = 0

for i, chunk in enumerate(pd.read_csv('big.csv',
                                       chunksize=chunksize)):
    rows_processed += len(chunk)
    # Report every 10 chunks
    if (i + 1) % 10 == 0:
        elapsed = time.time() - start
        rate = rows_processed / elapsed
        print(f'Chunk {i+1}: {rows_processed:,} rows '
              f'@ {rate/1000:.0f}k rows/sec')

print(f'Total: {rows_processed:,} rows in {time.time()-start:.1f}s')

Filtrando antes de agregar

Aplique filtros dentro de cada bloco antes de agregar, para evitar acumular dados indesejados. Por exemplo, se você se importar apenas com pedidos de 2024, filtre a coluna de data do bloco antes do groupby. Isso reduz a memória necessária para os resultados parciais e acelera a etapa final de concatenação. Sempre aplique os filtros o mais cedo possível no fluxo de processamento — um princípio fundamental do processamento eficiente de dados.

import pandas as pd

partials = []
for chunk in pd.read_csv('orders.csv',
                         chunksize=100000,
                         parse_dates=['order_date']):
    # Filter early: only 2024 orders
    mask = chunk['order_date'].dt.year == 2024
    filtered = chunk.loc[mask, ['category', 'revenue']]

    if len(filtered) > 0:
        p = filtered.groupby('category')['revenue'].sum()
        partials.append(p)

if partials:
    result = pd.concat(partials).groupby(level=0).sum()
    print(result)

Salvando resultados intermediários

Para tarefas de execução muito longa, salve periodicamente os resultados intermediários para poder retomar a partir de um ponto de verificação caso o processo seja interrompido. Grave as agregações de cada bloco em um arquivo Parquet ou CSV após cada N blocos. Se a tarefa falhar no bloco 800 de 1000, poderá recarregar as agregações salvas e continuar de onde parou, em vez de processar novamente o arquivo inteiro. Esse padrão de resiliência é essencial em pipelines de dados de produção.

import pandas as pd
import os

CHECKPOINT = 'checkpoint.csv'
running_total = 0.0
running_count = 0

# Resume from checkpoint if it exists
if os.path.exists(CHECKPOINT):
    ckpt = pd.read_csv(CHECKPOINT)
    running_total = ckpt['total'].iloc[0]
    running_count = int(ckpt['count'].iloc[0])
    print(f'Resuming from checkpoint: {running_count:,} rows')

for chunk in pd.read_csv('huge.csv', chunksize=100000):
    running_total += chunk['value'].sum()
    running_count += len(chunk)

# Save checkpoint
pd.DataFrame({'total': [running_total],
              'count': [running_count]}).to_csv(CHECKPOINT, index=False)
print(f'Final mean: {running_total / running_count:.4f}')

Verificação rápida

Teste sua compreensão dos conceitos de Análise de Dados desta lição.

Recapitulação da lição

Nesta lição, você aprendeu que: acumuladores em execução (sum, count, min/max, Series de frequências) permitem agregar arquivos grandes usando memória constante; o groupby em duas etapas (groupby parcial por bloco, seguido de concatenação e novo agrupamento) trata corretamente os grupos que atravessam vários blocos; e a filtragem antecipada em cada bloco reduz o custo da etapa de acumulação. A seguir, exploraremos Dask DataFrames como uma alternativa paralela e compatível diretamente com Pandas para grandes conjuntos de dados.

Perguntas Frequentes

A aula “Agregação incremental entre partes” é grátis?

Sim — o texto completo de “Agregação incremental entre partes” é grátis para ler aqui na web. Para praticá-la interativamente (um editor de código integrado e um tutor de IA 24/7) e desbloquear o restante do curso de Pandas & NumPy Academy, atualize para CoddyKit PRO. O curso de Pandas & NumPy Academy inclui 4 aulas no total.

O que vou aprender em “Agregação incremental entre partes”?

Acumule contagens, somas e valores mínimos/máximos em execução entre as partes sem armazenar o arquivo inteiro na memória. Você pratica Pandas & NumPy Academy com código prático que executa diretamente no navegador, e um tutor de IA 24/7 responde suas dúvidas enquanto trabalha na aula.

Preciso ter experiência prévia para começar Pandas & NumPy Academy?

Nenhuma experiência prévia é necessária. Pandas & NumPy Academy no CoddyKit é estruturado para alunos iniciantes até avançados, então você pode começar aqui ou desde o início e aprender no seu ritmo. Esta é a aula 2 de 4.

Quanto tempo leva a aula “Agregação incremental entre partes”?

A maioria das aulas CoddyKit leva cerca de 5–10 minutos. Cada uma é compacta e interativa, então você faz progresso constante e retoma exatamente de onde parou entre web e app.

Posso escrever e executar código nesta aula de Pandas & NumPy Academy?

Sim. Cada aula de Pandas & NumPy Academy inclui um editor de código integrado, então você escreve e executa código real direto no navegador e recebe feedback de IA instantaneamente — nenhuma configuração local necessária.

Todas as aulas deste curso

  1. CSV em fluxo com chunksize
  2. Agregação incremental entre partes
  3. Introdução aos DataFrames do Dask
  4. Parquet: armazenamento colunar rápido
← Voltar para Pandas & NumPy Academy