0Pricing
Pandas & NumPy Academy · Lección

Agregación incremental entre bloques

Acumule recuentos, sumas y valores mínimos y máximos entre bloques sin almacenar el archivo completo en memoria.

Agregación incremental entre bloques es una lección gratuita de Pandas & NumPy Academy en CoddyKit. Esta es la lección 2 de 4. Puedes leer la lección completa abajo gratuitamente — luego la practicas en el navegador con un editor de código integrado y un tutor de IA 24/7. Forma parte de la ruta de aprendizaje de Pandas & NumPy Academy, y tu progreso se sincroniza en la web y la app de CoddyKit. El curso de Pandas & NumPy Academy incluye 4 lecciones en total.

¿Por qué usar la agregación incremental?

La agregación incremental es la clave para analizar conjuntos de datos más grandes que la RAM sin distribuir el cálculo entre varias máquinas. En lugar de cargar todos los datos para calcular una estadística final, se mantienen acumuladores acumulados —sumas parciales, recuentos y valores mínimos/máximos— y se actualizan con cada bloque. El resultado final se ensambla a partir de estos acumuladores ligeros después de explorar todo el archivo. Este patrón permite procesar archivos de terabytes en un solo portátil.

Contar filas y calcular la media

Calcular la media entre bloques requiere realizar un seguimiento independiente de la suma y el recuento acumulados. No puede limitarse a promediar las medias de cada bloque, porque los bloques pueden tener tamaños diferentes. La fórmula correcta es total_sum / total_count. Este patrón se extiende a cualquier magnitud que pueda descomponerse: la varianza, la correlación y los histogramas también tienen fórmulas incrementales.

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 y máximo incrementales

Realizar un seguimiento del mínimo y el máximo globales entre bloques es sencillo: inicialice los valores con float('inf') y float('-inf') de Python y, después, actualícelos con el mínimo y el máximo de cada bloque. Esto evita almacenar datos intermedios. El patrón se generaliza al mínimo y máximo por grupo manteniendo un diccionario cuyas claves sean los identificadores de 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}')

Recuentos de frecuencia incrementales

Para las columnas categóricas, mantenga un diccionario de frecuencias acumuladas sumando el resultado de value_counts() de cada bloque a un acumulador de Series de Pandas. Como la suma de Series de Pandas alinea las etiquetas del índice, las categorías desconocidas de los bloques posteriores se incluyen automáticamente. Después de procesar todos los bloques, ordene por recuento para ver las categorías principales de todo el conjunto de datos.

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))

Agregación GroupBy incremental

Para calcular una suma o un recuento con groupby entre bloques, aplique groupby().agg() dentro de cada bloque y almacene la Series o el DataFrame resultante. Después del bucle, concatene todos los resultados parciales y aplique un segundo groupby para combinarlos. Este enfoque en dos etapas gestiona correctamente los grupos que aparecen en varios bloques, algo habitual cuando los datos están ordenados por fecha en lugar de 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))

Calcular la varianza de forma incremental (método de Welford)

Calcular la varianza entre bloques es más complicado que calcular la media. La fórmula ingenua E[X²] - E[X]² sufre cancelación catastrófica cuando las medias son grandes. El algoritmo en línea de Welford mantiene una media acumulada y una suma de desviaciones al cuadrado, y las actualiza con cada nuevo valor de forma numéricamente estable. Aunque SciPy lo implementa, comprender el patrón permite ampliarlo a la varianza ponderada y la covarianza.

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}')

Crear un histograma incremental

Calcular la distribución de una columna en un archivo demasiado grande para la RAM requiere un histograma incremental. Fije previamente los límites de los intervalos (basándose en una muestra pequeña o en el conocimiento del dominio), use después np.histogram(chunk_values, bins=edges) dentro de cada bloque y acumule los recuentos. Al final, represente los recuentos combinados en un gráfico de barras. Así es como sistemas de streaming como Kafka Streams y Flink calculan 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], '...')

Realizar un seguimiento aproximado de los valores únicos

Contar los valores distintos exactos entre bloques requiere almacenar todos los valores únicos, potencialmente millones. Para obtener recuentos aproximados a gran escala, use un esquema HyperLogLog, disponible en Python mediante la biblioteca hyperloglog. Como alternativa, realice un seguimiento de los valores únicos de cada bloque con un conjunto y calcule la unión, aunque esta crece sin límite. Para una aproximación económica, use pd.Series.nunique() por bloque e informe de la media: no es exacta, pero suele ser suficiente para perfilar los datos.

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

Informar del progreso durante ejecuciones largas

Procesar un archivo de varios gigabytes puede tardar varios minutos. Añada información sobre el progreso para saber que el pipeline sigue ejecutándose y poder estimar el tiempo restante. Cuente los bytes o las filas procesados y compárelos con el tamaño del archivo. La biblioteca tqdm hace que esto sea trivial con su envoltorio tqdm(reader). Incluso sin tqdm, imprimir una línea de estado cada 10 bloques proporciona información valiosa durante trabajos por lotes largos.

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')

Filtrar antes de agregar

Aplique los filtros dentro de cada bloque antes de realizar la agregación para evitar acumular datos no deseados. Por ejemplo, si solo le interesan los pedidos de 2024, filtre la columna de fecha del bloque antes de usar groupby. Esto reduce la memoria necesaria para los resultados parciales y acelera el paso final de concatenación. Aplique siempre los filtros lo antes posible en el pipeline: es un principio fundamental del procesamiento eficiente de datos.

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)

Guardar resultados intermedios

Para trabajos de muy larga duración, guarde periódicamente los resultados intermedios para poder reanudar desde un punto de control si el proceso se interrumpe. Escriba los agregados de cada bloque en un archivo Parquet o CSV después de cada N bloques. Si el trabajo falla en el bloque 800 de 1000, puede volver a cargar los agregados guardados y continuar desde donde lo dejó, en lugar de volver a procesar todo el archivo. Este patrón de resiliencia es esencial en las canalizaciones de datos de producción.

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}')

Comprobación rápida

Ponga a prueba su comprensión de los conceptos de análisis de datos de esta lección.

Repaso de la lección

En esta lección ha aprendido que los acumuladores en ejecución (suma, cantidad, mínimo/máximo y Series de frecuencias) permiten realizar agregaciones con un uso constante de memoria en archivos grandes; el groupby en dos etapas (un groupby parcial por bloque, seguido de concatenación y una nueva agrupación) gestiona correctamente los grupos distribuidos entre bloques; y el filtrado anticipado dentro de cada bloque reduce el coste de la etapa de acumulación. A continuación, exploraremos Dask DataFrames como sustituto paralelo de Pandas para conjuntos de datos grandes.

Preguntas frecuentes

¿La lección «Agregación incremental entre bloques» es gratis?

Sí — el texto completo de «Agregación incremental entre bloques» es gratis para leer aquí en la web. Para practicarla de forma interactiva (editor de código integrado y tutor de IA 24/7) y desbloquear el resto del curso de Pandas & NumPy Academy, actualiza a CoddyKit PRO. El curso de Pandas & NumPy Academy incluye 4 lecciones en total.

¿Qué aprenderé en «Agregación incremental entre bloques»?

Acumule recuentos, sumas y valores mínimos y máximos entre bloques sin almacenar el archivo completo en memoria. Practicas Pandas & NumPy Academy con código real que ejecutas directamente en el navegador, y un tutor de IA 24/7 responde tus preguntas mientras trabajas en la lección.

¿Necesito experiencia previa para empezar Pandas & NumPy Academy?

No se requiere experiencia previa. Pandas & NumPy Academy en CoddyKit está estructurado para principiantes hasta estudiantes avanzados, así que puedes empezar aquí o desde el inicio y avanzar a tu ritmo. Esta es la lección 2 de 4.

¿Cuánto tiempo toma la lección «Agregación incremental entre bloques»?

La mayoría de las lecciones de CoddyKit toman alrededor de 5–10 minutos. Cada una es compacta e interactiva, así que avanzas constantemente y retomas exactamente por donde dejaste en la web y la app.

¿Puedo escribir y ejecutar código en esta lección de Pandas & NumPy Academy?

Sí. Cada lección de Pandas & NumPy Academy incluye un editor de código integrado, así que escribes y ejecutas código real directamente en tu navegador y obtienes retroalimentación instantánea de IA — sin configuración local necesaria.

Todas las lecciones de este curso

  1. Transmitir CSV con chunksize
  2. Agregación incremental entre bloques
  3. Introducción a los DataFrames de Dask
  4. Parquet: almacenamiento columnar rápido
← Volver a Pandas & NumPy Academy