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