Procesamiento de consultas en Citus: arquitectura de ejecución distribuida

Un clúster de Citus consta de una instancia de coordinador y de varias instancias de trabajo. Los datos se particionan en los trabajos mientras que el coordinador almacena metadatos sobre estas particiones. Un cliente puede conectarse al coordinador o, en modo MX, a cualquier nodo de trabajo y enviar consultas allí. El nodo que recibe la consulta lo divide en fragmentos de consulta más pequeños donde cada fragmento de consulta se puede ejecutar de forma independiente en una partición. Ese nodo asigna los fragmentos de consulta a los trabajos que contienen las particiones pertinentes, supervisa su ejecución, combina sus resultados y devuelve el resultado final al cliente. En el diagrama siguiente se proporciona una breve descripción de la arquitectura de procesamiento de consultas.

Diagrama que muestra un clúster de Citus con un coordinador y nodos de trabajo, donde el nodo que recibe la consulta de cliente distribuye el trabajo a las particiones pertinentes.

La canalización de procesamiento de consultas de Citus implica dos componentes:

  • Distributed Query Planner y Executor
  • PostgreSQL Planner y Executor

En las secciones siguientes se describen estos componentes con mayor detalle.

Consulta de puntos de entrada y modo MX

En versiones anteriores de Citus, el coordinador era el único punto de entrada para las consultas de cliente. A partir de Citus 11, los clientes también pueden conectarse directamente a cualquier nodo de trabajo y enviar consultas allí. Esta funcionalidad se conoce como modo MX (extensión multiinquilino o masivamente paralela) y está habilitada de forma predeterminada en clústeres de Citus administrados.

El modo MX aborda dos problemas de escalado:

  • Detección activa del coordinador. Cuando cada cliente se conecta a un nodo, las ranuras de conexión, la CPU y la red de ese nodo se convierten en un cuello de botella. Permitir que los trabajadores acepten conexiones de cliente propagan esta carga en el clúster.
  • Latencia para las consultas locales de particiones. Cuando un cliente se conecta al trabajo que ya almacena la partición pertinente, Citus puede planear y ejecutar la consulta en la tabla de particiones local sin saltos de red adicionales.

Para que el modo MX funcione, cada nodo necesita una copia de up-tofecha de los metadatos del clúster (el catálogo de tablas distribuidas y ubicaciones de particiones). Citus sincroniza estos metadatos automáticamente con todos los nodos; la citus.metadata_sync_mode configuración controla cómo se ejecuta la sincronización.

Cuando un trabajador recibe una consulta, desempeña el mismo papel que desempeña el coordinador en la topología clásica:

  • Si la consulta toca particiones en otros trabajos, el nodo receptor abre conexiones a esos trabajos, distribuye los fragmentos y combina los resultados antes de devolverlos al cliente.
  • Si la consulta tiene como destino una partición que reside en el nodo receptor, el nodo usa citus.local_hostname para llegar a sí mismo y ejecuta el fragmento localmente. La optimización de planeación de rutas rápidas retrasadas (Citus 13.2) se basa en este caso para reutilizar los planes almacenados en caché y omitir el análisis, el análisis y los pasos del plan de la consulta de partición.

Las escrituras en tablas distribuidas y de referencia siguen la misma ruta de acceso: el nodo receptor aplica el protocolo de transacción distribuida estándar independientemente de si es el coordinador o un trabajo.

Planificador de consultas distribuidas

El planificador de consultas distribuidas de Citus toma una consulta SQL y la planea para la ejecución distribuida. El planificador se ejecuta en el nodo que recibió la consulta: el coordinador de la topología clásica o cualquier trabajo en modo MX.

En SELECT el caso de las consultas, el planificador crea primero un árbol de plan de la consulta de entrada y lo transforma en su formulario asociativo y conmutante para que se pueda paralelizar. También aplica varias optimizaciones para asegurarse de que las consultas se ejecutan de forma escalable y que se minimiza la E/S de red.

A continuación, el planificador divide la consulta en dos partes: la consulta de nivel superior, que se ejecuta en el nodo receptor y los fragmentos de consulta de trabajo, que se ejecutan en particiones individuales en los trabajos. A continuación, el planificador asigna los fragmentos de consulta a los trabajos de forma que todos sus recursos se usen de forma eficaz. Después de este paso, el plan de consulta distribuida se pasa al ejecutor distribuido para su ejecución.

El proceso de planeación de las búsquedas de clave-valor en la columna de distribución o las consultas de modificación es ligeramente diferente, ya que alcanzan exactamente una partición. Cuando el planificador recibe una consulta entrante, decide la partición correcta a la que se debe enrutar la consulta. Determina la partición correcta para la consulta mediante la extracción de la columna de distribución en la fila entrante y la búsqueda de los metadatos. A continuación, el planificador vuelve a escribir el CÓDIGO SQL de ese comando para hacer referencia a la tabla de particiones en lugar de a la tabla original. A continuación, este plan reescrito se pasa al ejecutor distribuido.

Planificación de rutas rápidas retrasadas (Citus 13.2)

En Citus 13.2, el planificador retrasa la creación del plan de marcador de posición de ruta rápida hasta que identifica la partición. En el modo MX, si la selección de ubicación de particiones es local para el nodo que controla la consulta de cliente, Citus puede evitar el análisis, el análisis y los pasos del plan para la consulta de particiones y reutilizar un plan almacenado en caché. Este enfoque mejora el rendimiento.

Elegibilidad:

  • La consulta es SELECT o UPDATE en una tabla distribuida (esquema o particionado de columna) o en una tabla local administrada por Citus.
  • No hay funciones volátiles.
  • La partición se puede determinar en el momento del plan y es local para el nodo receptor.
  • Actualmente no se admiten tablas de referencia.

Comportamiento:

  • Si la partición es local y segura, el ejecutor reemplaza el OID de la tabla distribuida por el OID de partición, llama a standard_plannery almacena en caché el plan en la tarea.
  • De lo contrario, el ejecutor vuelve al plan de marcador de posición de ruta de acceso rápido.

GUC:

  • citus.enable_local_fast_path_query_optimization (valor predeterminado on).

Ejecutor de consultas distribuidas

El ejecutor distribuido de Citus ejecuta planes de consulta distribuidos y controla los errores. El ejecutor es adecuado para obtener respuestas rápidas a las consultas que implican filtros, agregaciones y combinaciones coubicadas. También es bueno para ejecutar consultas de un solo inquilino con cobertura sql completa. El ejecutor se ejecuta en el nodo receptor y abre una conexión por partición a los trabajos que contienen esas particiones, enviando todas las consultas de fragmentos a ellas. En el modo MX, el nodo receptor podría ser un trabajo y las particiones que almacena localmente se ejecutan sin un salto de red. A continuación, el ejecutor captura los resultados de cada consulta de fragmento, los combina y devuelve los resultados finales al cliente.

Ejecución de inserción y extracción de subconsultas y CTE

Si es necesario, Citus puede recopilar resultados de subconsultas y expresiones de tabla comunes (CTE) en el nodo receptor y, a continuación, insertarlos en los trabajos para que los use una consulta externa. Esta arquitectura permite a Citus admitir una mayor variedad de construcciones SQL.

Por ejemplo, tener subconsultas en una cláusula WHERE no siempre se puede ejecutar en línea al mismo tiempo que la consulta principal, pero debe realizarse por separado. Supongamos que una aplicación de análisis web mantiene una page_views tabla particionada por page_id. Para consultar el número de hosts de visitante en las 20 páginas más visitadas, use una subconsulta para buscar la lista de páginas y, a continuación, una consulta externa para contar los hosts.

SELECT page_id, count(distinct host_ip)
FROM page_views
WHERE page_id IN (
  SELECT page_id
  FROM page_views
  GROUP BY page_id
  ORDER BY count(*) DESC
  LIMIT 20
)
GROUP BY page_id;

El ejecutor ejecuta un fragmento de esta consulta en cada partición por page_id, cuenta s distintos host_ipy combina los resultados en el coordinador. Sin embargo, en LIMIT la subconsulta significa que la subconsulta no se puede ejecutar como parte del fragmento. Al planear de forma recursiva la consulta Citus puede ejecutar la subconsulta por separado, insertar los resultados en todos los trabajos, ejecutar la consulta de fragmento principal y devolver los resultados al coordinador. El diseño de inserción y extracción admite subconsultas como la del ejemplo anterior.

Para ver esta ejecución de inserción y extracción en acción, revise la salida EXPLAIN de esta consulta.

GroupAggregate (cost=0.00..0.00 rows=0 width=0)
  Group Key: remote_scan.page_id
  -> Sort (cost=0.00..0.00 rows=0 width=0)
    Sort Key: remote_scan.page_id
    -> Custom Scan (Citus Adaptive) (cost=0.00..0.00 rows=0 width=0)
      -> Distributed Subplan 6_1
        -> Limit (cost=0.00..0.00 rows=0 width=0)
          -> Sort (cost=0.00..0.00 rows=0 width=0)
            Sort Key: COALESCE((pg_catalog.sum((COALESCE((pg_catalog.sum(remote_scan.worker_column_2))::bigint, '0'::bigint))))::bigint, '0'::bigint) DESC
            -> HashAggregate (cost=0.00..0.00 rows=0 width=0)
              Group Key: remote_scan.page_id
              -> Custom Scan (Citus Adaptive) (cost=0.00..0.00 rows=0 width=0)
                Task Count: 32
                Tasks Shown: One of 32
                -> Task
                  Node: host=localhost port=9701 dbname=postgres
                  -> HashAggregate (cost=54.70..56.70 rows=200 width=12)
                    Group Key: page_id
                    -> Seq Scan on page_views_102008 page_views (cost=0.00..43.47 rows=2247 width=4)
      Task Count: 32
      Tasks Shown: One of 32
      -> Task
        Node: host=localhost port=9701 dbname=postgres
        -> HashAggregate (cost=84.50..86.75 rows=225 width=36)
          Group Key: page_views.page_id, page_views.host_ip
          -> Hash Join (cost=17.00..78.88 rows=1124 width=36)
            Hash Cond: (page_views.page_id = intermediate_result.page_id)
            -> Seq Scan on page_views_102008 page_views (cost=0.00..43.47 rows=2247 width=36)
            -> Hash (cost=14.50..14.50 rows=200 width=4)
              -> HashAggregate (cost=12.50..14.50 rows=200 width=4)
                Group Key: intermediate_result.page_id
                -> Function Scan on read_intermediate_result intermediate_result (cost=0.00..10.00 rows=1000 width=4)

El proceso está bastante implicado, así que vamos a separarlo y examinar cada pieza.

GroupAggregate (cost=0.00..0.00 rows=0 width=0)
  Group Key: remote_scan.page_id
  -> Sort (cost=0.00..0.00 rows=0 width=0)
    Sort Key: remote_scan.page_id

La raíz del árbol es lo que hace el nodo de coordinación con los resultados de los trabajos. En este caso, los agrupa y GroupAggregate requiere que se ordenen primero.

    -> Custom Scan (Citus Adaptive) (cost=0.00..0.00 rows=0 width=0)
      -> Distributed Subplan 6_1

El examen personalizado tiene dos subárboles grandes, empezando por un subplan distribuido.

        -> Limit (cost=0.00..0.00 rows=0 width=0)
          -> Sort (cost=0.00..0.00 rows=0 width=0)
            Sort Key: COALESCE((pg_catalog.sum((COALESCE((pg_catalog.sum(remote_scan.worker_column_2))::bigint, '0'::bigint))))::bigint, '0'::bigint) DESC
            -> HashAggregate (cost=0.00..0.00 rows=0 width=0)
              Group Key: remote_scan.page_id
              -> Custom Scan (Citus Adaptive) (cost=0.00..0.00 rows=0 width=0)
                Task Count: 32
                Tasks Shown: One of 32
                -> Task
                  Node: host=localhost port=9701 dbname=postgres
                  -> HashAggregate (cost=54.70..56.70 rows=200 width=12)
                    Group Key: page_id
                    -> Seq Scan on page_views_102008 page_views (cost=0.00..43.47 rows=2247 width=4)

Los nodos de trabajo ejecutan este subplan para cada una de las 32 particiones (Citus elige un representante para mostrar). Puede reconocer todas las partes de la IN (...) subconsulta: la ordenación, la agrupación y la limitación. Cuando todos los trabajos completan esta consulta, devuelven su salida al coordinador, lo que lo coloca como resultados intermedios.

      Task Count: 32
      Tasks Shown: One of 32
      -> Task
        Node: host=localhost port=9701 dbname=postgres
        -> HashAggregate (cost=84.50..86.75 rows=225 width=36)
          Group Key: page_views.page_id, page_views.host_ip
          -> Hash Join (cost=17.00..78.88 rows=1124 width=36)
            Hash Cond: (page_views.page_id = intermediate_result.page_id)

Citus inicia otro trabajo del ejecutor en este segundo subárbol. Va a contar hosts distintos en page_views. Usa join para conectarse con los resultados intermedios. Los resultados intermedios ayudan a restringirlo a las 20 páginas principales.

            -> Seq Scan on page_views_102008 page_views (cost=0.00..43.47 rows=2247 width=36)
            -> Hash (cost=14.50..14.50 rows=200 width=4)
              -> HashAggregate (cost=12.50..14.50 rows=200 width=4)
                Group Key: intermediate_result.page_id
                -> Function Scan on read_intermediate_result intermediate_result (cost=0.00..10.00 rows=1000 width=4)

El trabajo recupera internamente los resultados intermedios mediante una read_intermediate_result función , que carga datos de un archivo en el que se copió el nodo de coordinación.

En este ejemplo se muestra cómo Citus ejecutó la consulta en varios pasos con un subplan distribuido y cómo puede usar EXPLAIN para obtener información sobre la ejecución de consultas distribuidas.

Organizador y ejecutor de PostgreSQL

Después de que el ejecutor distribuido envíe los fragmentos de consulta, cada trabajo procesa sus fragmentos como las consultas normales de PostgreSQL. En el modo MX, el trabajo receptor también ejecuta fragmentos para las particiones que almacena localmente. PostgreSQL Planner en cada trabajo elige el plan más óptimo para ejecutar la consulta localmente en la tabla de particiones correspondiente. El ejecutor de PostgreSQL ejecuta la consulta y devuelve los resultados de la consulta al ejecutor distribuido. Para obtener más información sobre postgreSQL Planner y executor, consulte el manual de PostgreSQL. Por último, el ejecutor distribuido pasa los resultados al nodo receptor para la agregación final y los devuelve al cliente.