Por qué COUNT(DISTINCT) desactiva el paralelismo en Postgres

Fuentes: The DISTINCT in your COUNT

En Postgres, una consulta tan habitual como SELECT count(DISTINCT user_id) FROM events no puede aprovechar los workers paralelos, ni aunque la máquina disponga de varios núcleos y la configuración lo permita. El motivo está en el modelo de ejecución del agregador, no en una estimación de coste ni en un índice mal elegido.

El artículo lo demuestra con un experimento reproducible sobre una tabla events con diez millones de filas y unos 50.000 usuarios distintos. Con count(*) el planificador lanza cuatro workers en un Parallel Seq Scan, cada uno calcula un Partial Aggregate y el líder suma los cinco recuentos parciales. Con count(DISTINCT user_id), basta añadir esa palabra para que el plan colapse a un Aggregate serial sobre un Sort que se desborda a disco (115 MB de temp), ejecutado por un único proceso.

La razón es que las agregaciones en paralelo funcionan en dos fases: cada worker construye un estado parcial mediante una función combine que el agregador sabe fusionar. Para count, sum, avg, min o max esa función existe. Para count(DISTINCT x) no existe una forma de combinar dos recuentos parciales que sea correcta sin intercambiar el conjunto completo de valores vistos: habría que enviar a un único nodo todos los usuarios distintos de todos los workers, anulando el beneficio del paralelismo. Sin Partial Aggregate válido, el Gather no tiene qué alimentar y el escaneo paralelo deja de tener sentido. La cláusula debug_parallel_query confirma el diagnóstico: fuerza un Gather con Single Copy: true que solo reencamina la salida, sin repartir trabajo real.

El problema se extiende a toda la lista SELECT: un sum(amount) perfectamente paralelizable, si comparte bloque de consulta con count(DISTINCT user_id), acaba ejecutándose en serie bajo el mismo nodo Aggregate. FILTER no presenta este inconveniente y mantiene el plan paralelo.

La solución pasa por reescribir la deduplicación como un GROUP BY, que sí tiene modo parcial: cada worker construye un hash parcial de grupos, el líder los fusiona con HashAggregate y el conteo final es trivial. Con la reescritura propuesta, la misma consulta vuelve a ejecutarse con cuatro workers activos y un consumo de memoria de unos 3 MB.