Menu fechado

Abstração de Grafos de Alta Performance na Netflix: Arquitetura e Escala

Abstração de Grafos

🚀 Introdução à Abstração de Grafos na Netflix

A Netflix opera com um ecossistema vasto e complexo, onde o uso de grafos é fundamental para sustentar diversas necessidades de negócio. Para otimizar a performance e a escalabilidade, a empresa categoriza seus casos de uso de grafos em duas vertentes principais:

OLAP (Online Analytical Processing)

Focado na exploração algorítmica de grandes conjuntos de dados, o OLAP prioriza a análise profunda em vez da velocidade imediata. Estes casos utilizam modelos e linguagens padrão de mercado, como RDF com SPARQL, Property Graphs com Gremlin ou openCypher, e até mesmo SQL. O objetivo aqui não é a baixa latência, mas sim a capacidade de processamento analítico complexo, frequentemente envolvendo varreduras completas (full-graph scans) e agregações de alto custo computacional.

OLTP (Online Transactional Processing)

Diferente do OLAP, os casos de uso OLTP exigem uma performance extrema: a capacidade de processar milhões de operações por segundo com resultados de travessia entregues em milissegundos. Para atingir esse nível de eficiência, a Netflix adota compromissos arquiteturais estratégicos, como a aceitação de consistência eventual e a imposição de restrições na complexidade das consultas (como profundidade máxima de travessia e limites de fan-out).

A Abstração de Grafos da Netflix foi projetada especificamente para atender a essa segunda categoria. Atualmente, essa infraestrutura processa cerca de 10 milhões de operações por segundo em 650 TB de dados, garantindo alta disponibilidade global e eficiência de custos, sendo essencial para experiências em tempo real e sistemas de streaming de alta performance.

🏗️ Arquitetura e Modelo de Dados

Em vez de desenvolver camadas de persistência e cache do zero, a Abstração de Grafos da Netflix foi construída sobre abstrações de dados já consolidadas na infraestrutura da empresa, garantindo escalabilidade e eficiência operacional através de componentes de infraestrutura resilientes.

Integração com o Ecossistema de Abstrações

A arquitetura utiliza a Key-Value (KV) Abstraction como o alicerce para o armazenamento da visão atual de nós e arestas, servindo como o índice de tempo real para todas as consultas. Para cenários que exigem análise histórica da evolução do grafo, o sistema permite a integração opcional com a TimeSeries (TS) Abstraction. Para garantir latências na casa dos milissegundos, a solução utiliza o EVCache, um datastore distribuído em memória baseado em Memcached, enquanto o Data Gateway Control Plane atua como a camada de controle centralizada para gerenciar esquemas, orquestrar metadados e automatizar o provisionamento e a configuração dos datasets em clusters multi-região.

Modelo Property Graph

A abstração adota o modelo Property Graph, onde o grafo é composto por nós e arestas de diversos tipos, cada um contendo propriedades fortemente tipadas. Essa tipagem rigorosa é fundamental para permitir filtragens eficientes e garantir a consistência em exportações de dados para pipelines de processamento em lote (batch processing). As arestas podem ser configuradas como unidirecionais ou bidirecionais, dependendo da semântica da relação e dos requisitos de travessia.

Namespaces e Isolamento

Para manter a organização e a performance, os dados são segregados em unidades lógicas chamadas “namespaces”. Cada namespace é mapeado para uma camada de armazenamento físico via Data Gateway Control Plane, podendo ser implantado em hardware dedicado ou compartilhado. A automação de provisionamento da Netflix determina a configuração de hardware mais econômica com base em requisitos críticos, como throughput, latência de p99 e tamanho do dataset, garantindo o isolamento de falhas (blast radius reduction).

Gestão de Esquemas (Graph Schema)

Cada namespace possui um esquema explícito que define tipos de nós, arestas, propriedades permitidas e regras de relacionamento. O esquema é implementado como uma coleção de mapeamentos de arestas (edge mappings). Por exemplo:


{
"edgeConfig": {
"edgeMappings": [
{
"edgeMappingKey": { "fromNodeType": "account", "edgeType": "owns", "toNodeType": "profile" },
"directionType": "UNIDIRECTIONAL"
}
]
}
}

Ao carregar esse esquema em tempo de execução, os servidores da Abstração constroem um grafo de metadados em memória, o que possibilita otimizações críticas:

  • Qualidade de Dados: Rejeição automática de escritas que não conformam com o esquema, prevenindo a corrupção do grafo.
  • Planejamento de Consultas: Construção rápida de caminhos de travessia otimizados via otimizadores baseados em custo (CBO).
  • Deduplicação: Eliminação de processamento redundante em travessias bidirecionais através de lógica de normalização.
  • Eliminação de Caminhos: Remoção de rotas impossíveis ou incompatíveis com os filtros da consulta, reduzindo o espaço de busca.

🔍 O Papel do Esquema de Grafo

O esquema de grafo atua como a espinha dorsal da Abstração de Grafos, sendo carregado em memória pelos servidores durante a inicialização para formar um grafo de metadados das relações possíveis. Esta estrutura declarativa permite que o sistema execute otimizações fundamentais que garantem tanto a integridade do sistema quanto a performance em escala.

Garantia de Qualidade de Dados

O esquema impõe um contrato rigoroso sobre a topologia do grafo. Durante as operações de escrita, a Abstração valida automaticamente se os nós, arestas e propriedades estão em conformidade com as definições pré-estabelecidas no Data Gateway Control Plane. Qualquer tentativa de persistir dados que violem os tipos de propriedades, direções de arestas ou relacionamentos proibidos é rejeitada, garantindo que o dataset permaneça consistente e pronto para exportações confiáveis.

Planejamento de Consultas (Query Planning)

Com o conhecimento prévio dos tipos de nós e das arestas permitidas, o motor de travessia utiliza o esquema para construir caminhos de execução otimizados. Ao receber uma consulta, o sistema mapeia as rotas possíveis entre os nós, permitindo que o planejador de consultas identifique a sequência mais eficiente de saltos (hops) para alcançar o resultado desejado, minimizando o processamento desnecessário e o uso de recursos de CPU.

Eliminação de Caminhos de Travessia

A consciência do esquema permite que a Abstração reduza drasticamente o espaço de busca durante a execução de consultas complexas:

  • Filtragem Semântica: O sistema remove automaticamente caminhos que envolvem relacionamentos impossíveis ou incompatíveis com os tipos de nós solicitados.
  • Deduplicação de Arestas: Em travessias bidirecionais entre o mesmo tipo de nó, o esquema permite identificar e eliminar caminhos redundantes, evitando que o mesmo processamento ocorra múltiplas vezes.
  • Otimização de Filtros: Se uma consulta inclui filtros em propriedades, o esquema permite descartar antecipadamente qualquer caminho que não suporte as chaves ou tipos de dados especificados, reduzindo a carga de I/O e o overhead de serialização.

Periodicamente, os servidores atualizam seu grafo de metadados em memória consultando o Control Plane, garantindo que qualquer alteração no esquema seja propagada sem a necessidade de reinicialização, mantendo a flexibilidade operacional enquanto preserva a performance otimizada.

💾 Indexação em Tempo Real e Armazenamento

A Abstração de Grafos da Netflix utiliza a Key-Value (KV) Abstraction como base para seu armazenamento persistente, garantindo alta disponibilidade e latência na casa dos milissegundos. A estratégia de indexação é fundamentada em uma estrutura de “mapa de mapas ordenados”, que oferece flexibilidade para diversos padrões de acesso e alta densidade de armazenamento.

Armazenamento de Nós

Cada tipo de nó é isolado em seu próprio namespace KV, onde todas as propriedades associadas são armazenadas. Essa abordagem permite:




  • Leituras eficientes: O acesso a um nó e todas as suas propriedades ocorre em uma única busca de partição (point lookup).
  • Pushdown de operações: A seleção e filtragem de propriedades são delegadas à camada KV, reduzindo drasticamente o volume de dados trafegados na rede.
  • Exportação paralela: A separação por tipos de nós facilita a exportação de dados em larga escala para data lakes.

Armazenamento de Arestas: Links vs. Propriedades

Para arestas, o sistema adota uma separação estratégica entre os links (conexões entre nós) e as propriedades das arestas. Esta decisão arquitetural traz benefícios críticos:

  • Prevenção de “Wide Rows”: Ao desacoplar os links das propriedades, evitamos a criação de partições gigantescas em bancos de dados como o Cassandra, permitindo que o sistema suporte milhões de conexões por nó sem degradação de performance.
  • Upserts otimizados: É possível atualizar propriedades de uma aresta individualmente sem a necessidade de ler ou reescrever todo o conjunto de dados daquela conexão.

Como contrapartida, essa separação exige que as escritas não sejam atômicas entre namespaces, um desafio mitigado pela nossa estratégia de consistência eventual e reconciliação assíncrona.

Índices Forward e Reverse

Para viabilizar travessias bidirecionais eficientes, a Abstração mantém índices separados para as direções forward (origem para destino) e reverse (destino para origem).

Para garantir a integridade dos identificadores durante mutações, utilizamos um identificador agnóstico à direção, gerado através da concatenação lexicográfica dos IDs dos nós de origem e destino. Isso permite que qualquer propriedade de aresta seja acessada ou modificada em uma única chamada de banco de dados, independentemente da direção da travessia solicitada. O sistema suporta:

  • Point Reads: Recuperação de propriedades via ID da aresta em uma única busca.
  • Range Reads: Varredura eficiente de vizinhos a partir de um nó de origem, utilizando o índice de links apropriado.
  • Ordenação: Embora o padrão seja a ordenação lexicográfica, o sistema pode realizar ordenações em memória baseadas no timestamp de última escrita para recuperar as conexões mais recentes.

⚡ Estratégias de Cache e Performance

Embora a camada de armazenamento persistente ofereça alta disponibilidade, a natureza altamente interconectada dos grafos impõe desafios severos de performance. A Abstração de Grafos da Netflix implementa estratégias de cache multicamadas para mitigar dois fenômenos críticos: a amplificação de escrita e a amplificação de leitura.

Mitigação da Amplificação de Escrita: Cache Write-Aside para Links

Cada operação de escrita pode disparar múltiplas atualizações em índices distintos, gerando uma carga desnecessária no armazenamento durável. Para otimizar esse processo, utilizamos uma estratégia de write-aside especificamente para os links de arestas (edge links). Como um link de aresta contém apenas a conexão em si e um timestamp de última escrita, o cache atua como um filtro de redundância:

  • Prevenção de escritas redundantes: O sistema verifica o cache antes de persistir um link. Se o link já existir, a escrita no armazenamento durável é evitada.
  • Controle de consistência: Este mecanismo é equilibrado por janelas de TTL (Time-To-Live) configuráveis, invalidação explícita em operações de exclusão e aquisição de leases com exponential backoff, garantindo que o timestamp de última escrita permaneça atualizado dentro dos limites de tolerância de estagnação.

Mitigação da Amplificação de Leitura: Cache Read-Aside para Propriedades

Uma única requisição de travessia pode resultar em milhares de operações de busca no backend. Para reduzir essa carga, a Abstração integra-se ao EVCache, o datastore distribuído em memória da Netflix, aplicando uma estratégia de read-aside:

  • Eficiência de Custo: Múltiplos namespaces de KV (Key-Value) compartilham os mesmos clusters de cache, otimizando a utilização de recursos de memória RAM.
  • Granularidade: O cache é aplicado tanto no nível de registro quanto no nível de item, permitindo que leituras subsequentes sejam servidas diretamente da memória com latência de sub-milissegundos.

Estratégias de Invalidação

A escolha da estratégia de invalidação no EVCache é determinada pelos requisitos específicos de throughput e consistência de cada caso de uso:

  • Invalidation on Write: Ideal para grafos com baixa frequência de alteração que exigem consistência estrita. Cada escrita invalida os caches de registro e item, garantindo que não haja dados obsoletos, ao custo de um maior tráfego no cache.
  • TTL-driven Invalidation: Recomendada para objetos modificados frequentemente. A expiração baseada em tempo permite uma maior tolerância à estagnação, reduzindo a pressão sobre o sistema de invalidação.

Atualmente, a equipe está evoluindo para uma arquitetura de Write-Through Caching, que visa armazenar a maioria dos dados necessários durante as travessias, permitindo a organização de índices por diferentes ordens de classificação (como o timestamp de última escrita) em troca de um consumo otimizado de memória.

⚖️ Consistência e Travessias

Garantir a integridade dos dados em um sistema distribuído de alta escala impõe desafios significativos. A natureza interconectada dos grafos, somada à exigência de latência na casa dos milissegundos, levou a Netflix a adotar um modelo de consistência eventual estrita entre múltiplas regiões.

Gerenciamento de Consistência

Para manter o sistema íntegro, a Abstração de Grafos implementa estratégias específicas:

  • Reparo de Entropia: Como cada operação de escrita persiste dados em múltiplos índices e namespaces de forma paralela para maximizar o throughput, o sistema utiliza um mecanismo robusto de retentativa via Kafka para mitigar falhas parciais e prevenir a divergência de dados.
  • Exclusão Assíncrona de Nós: A remoção de um nó em um grafo altamente conectado não pode ser síncrona, dado o volume de arestas associadas. O processo de exclusão ocorre de forma assíncrona para não impactar a latência das chamadas. Para garantir a correção durante atualizações concorrentes, utiliza-se o mecanismo de resolução de conflitos Last-Write-Wins (LWW).
  • Replicação Global: Tanto a camada de cache quanto o armazenamento durável replicam dados de forma assíncrona entre regiões, assegurando que o sistema permaneça disponível globalmente, mesmo sob falhas regionais ou particionamento de rede.

API de Travessia via gRPC

A Abstração expõe uma API de travessia customizada baseada em gRPC, inspirada na linguagem Gremlin. Esta interface permite que os desenvolvedores naveguem pelo grafo distribuído de maneira eficiente, suportando:

  • Encadeamento de Consultas: Navegação entre nós através de múltiplos saltos (hops).
  • Filtros e Ordenação: Aplicação de critérios de filtragem por tipo de nó/aresta e ordenação de resultados.
  • Controle de Fanout: Limitação de resultados para evitar sobrecarga na rede e processamento.

Um exemplo de requisição de travessia, estruturada via TraversalRequest.newBuilder(), permite definir o ponto de partida, a direção da navegação (IN ou OUT) e a seleção específica de propriedades, otimizando a troca de dados entre o serviço e a camada de armazenamento.

📈 Resultados e Conclusão

A arquitetura de Abstração de Grafos da Netflix demonstra robustez operacional ao suportar um volume massivo de 10 milhões de operações por segundo, distribuídas em 650 TB de dados. A eficiência do sistema é evidenciada pela baixa latência alcançada em cenários críticos de produção:

Desempenho em Tempo Real

  • Persistência: Tanto a leitura quanto a escrita de nós e arestas operam com latências de milissegundos de um dígito (p99).
  • Travessias: Consultas de 1 salto são executadas consistentemente na casa dos milissegundos de um dígito.
  • Escalabilidade de Consultas: Mesmo em cenários complexos, como o Real-Time Distributed Graph (RDG) que utiliza travessias de 2 saltos com alto fan-out, o sistema mantém o p90 abaixo de 50ms.
  • Operações Assíncronas: Processos de manutenção, como a exclusão de nós, são realizados com latência sub-segundo, garantindo que a integridade do grafo não comprometa a experiência em tempo real.

A Abstração de Grafos consolidou-se como um pilar fundamental para a infraestrutura de dados da Netflix. Ao oferecer uma solução que equilibra alta disponibilidade, custo-eficiência e performance, a plataforma permite que a engenharia foque na criação de novas funcionalidades em vez de gerenciar a complexidade do armazenamento subjacente. À medida que a Netflix expande sua presença em verticais como jogos, anúncios e conteúdo ao vivo, esta arquitetura continuará sendo o motor essencial para extrair valor de conexões complexas, garantindo que a plataforma permaneça ágil e escalável frente aos desafios futuros.


Fonte: netflixtechblog.com
Curadoria e Insights: Redação YTI&W (Developers).



Deixe um comentário

O seu endereço de e-mail não será publicado. Campos obrigatórios são marcados com *

Publicado em:Banco de Dados,Engenharia de Dados,Engenharia de Software,Performance & Otimização
Fale Conosco
×

Inscreva-se em nossa Newsletter!


Receba nossos lançamentos e artigos em primera mão!