Este repositório contém dois dos três componentes do sistema distribuído de séries temporais:
- AST (Armazenamento de Séries Temporais) — armazena eventos em CSV por dia, serve pedidos de leitura via gRPC e forma uma DHT com consistent hashing
- SA (Servidor de Agregação) — recebe pedidos de agregação, obtém eventos do AST e calcula resultados
A comunicação entre SA e AST usa gRPC com RxJava (rx3grpc) e serialização protobuf.
SSSimulator (Dealer) ──ZMQ──► SARouter → AggregationService → ASTClient ──gRPC──► FetchEventsHandler
│
ASTService
(DHT routing)
│
EventStorage (CSV)
- Java 17+
- Maven 3.8+
- protoc instalado e no PATH
Verificar instalações:
java -version
mvn -version
protoc --versionJava_Modules/
├── pom.xml
├── config_ast1.properties # Configuração do nó AST-1
├── config_ast2.properties # Configuração do nó AST-2
├── config_ast3.properties # Configuração do nó AST-3
├── config_sa.properties # Configuração do SA
└── src/
├── main/java/
│ ├── Common/
│ │ ├── model/
│ │ │ ├── Event.java # Modelo de domínio central
│ │ │ ├── AggregationQuery.java # Pedido de agregação (partilhado SA+AST)
│ │ │ └── AggregationResult.java # Resultado de uma agregação
│ │ └── util/
│ │ ├── EventConverter.java # Event Java ↔ EventProto
│ │ └── QueryConverter.java # AggregationQuery Java ↔ proto
│ │
│ ├── AST/
│ │ ├── ASTMain.java # Entry point do AST
│ │ ├── ASTService.java # Lógica de leitura, filtragem e DHT routing
│ │ ├── EventStorage.java # Persistência em CSV
│ │ ├── dht/
│ │ │ ├── ConsistentHashRing.java # Anel de consistent hashing
│ │ │ ├── DHTManager.java # Gestão do anel e routing entre nós
│ │ │ └── ASTNode.java # Representação de um nó AST
│ │ ├── grpc/
│ │ │ ├── ASTGrpcServer.java # Servidor gRPC (lifecycle)
│ │ │ └── FetchEventsHandler.java # Handler dos pedidos gRPC (fetchEvents + storeEvent)
│ │ └── util/
│ │ └── TestDataGenerator.java # Gerador de dados de teste com suporte DHT
│ │
│ └── AS/
│ ├── SAMain.java # Entry point do SA
| ├── AggregationService.java # Orquestrador (cache + AST + engine)
│ ├── AggregationCache.java # Cache de resultados em memória
| ├── AggregationService.java # Orquestrador (cache + AST + engine)
| └── AggregationEngine.java # Motor de cálculo (COUNT, SUM, etc.)
│ ├── util/
│ │ └── SSSimulator.java # Simulador do SS para testes
│ │ └── ZMQTestClient.java # Cliente ZMQ para testes manuais
│ └── grpc/
│ └── ASTClient.java # Cliente gRPC do AST (com nó aleatório)
│ └── zmq/
│ └── SARouter.java # ROUTER ZMQ que trata de pedidos do SS
│
└── test/java/
└── AST/dht/
└── ConsistentHashRingTest.java # Testes unitários do consistent hashing
ast.nodes=localhost:50051,localhost:50052,localhost:50053
ast.self=localhost:50051 # endereço deste nó
ast.port=50051 # porta deste nó
ast.vnodes=150 # número de virtual nodessa.ast.nodes=localhost:50051,localhost:50052,localhost:50053
sa.port=6000Na raiz de Java_Modules/:
mvn compileEm três janelas de terminal separadas:
# Janela 1
mvn exec:java "-Dexec.mainClass=AST.ASTMain" "-Dexec.args=config_ast1.properties"
# Janela 2
mvn exec:java "-Dexec.mainClass=AST.ASTMain" "-Dexec.args=config_ast2.properties"
# Janela 3
mvn exec:java "-Dexec.mainClass=AST.ASTMain" "-Dexec.args=config_ast3.properties"Output esperado em cada janela:
[AST-gRPC] Servidor iniciado na porta 5005X
mvn exec:java "-Dexec.mainClass=AST.util.TestDataGenerator"Gera 50 eventos por dia durante 7 dias (2025-01-01 a 2025-01-07) e envia cada dia para o nó AST correcto via DHT. Output esperado:
Dia 2025-01-01 → nó localhost:50053 (50 eventos)
Dia 2025-01-02 → nó localhost:50053 (50 eventos)
...
mvn exec:java "-Dexec.mainClass=AS.SAMain" "-Dexec.args=config_sa.properties"Output esperado:
[SA] Started Aggregation Server
--- COUNT alarmes ---
Cache miss: AggregationQuery{op=COUNT, days=[2025-01-01..2025-01-07], type==alarme}
AggregationResult{value=721.0}
--- SUM speed em alarmes ---
Cache miss: AggregationQuery{op=SUM, ..., k2=speed}
AggregationResult{value=79304.1}
--- SUM_OF_PRODUCTS speed*weight em travagens ---
Cache miss: AggregationQuery{op=SUM_OF_PRODUCTS, ..., k2=speed, k3=weight}
AggregationResult{value=4.13E7}
--- COUNT alarmes (cache hit) ---
Cache hit: AggregationQuery{op=COUNT, ...}
AggregationResult{value=721.0} ← mesmo valor que o primeiro COUNT
[SA-ZMQ] ROUTER bound on port 5555
Numa janela separada:
mvn exec:java "-Dexec.mainClass=AS.util.ZMQTestClient"Output esperado:
req-001 → {"request_id":"req-001","status":"OK","value":109.0}
req-002 → {"request_id":"req-002","status":"OK","value":12561.43}
req-003 → {"request_id":"req-003","status":"OK","value":109.0}
Output esperado no SA:
Cache hit: AggregationQuery{op=COUNT, days=[2025-01-01..2025-01-07], type==alarme}
Cache hit: AggregationQuery{op=SUM, days=[2025-01-01..2025-01-07], type==alarme, k2=speed}
Cache hit: AggregationQuery{op=COUNT, days=[2025-01-01..2025-01-07], type==alarme}
| Componente | O que valida |
|---|---|
ConsistentHashRing |
Determinismo, distribuição uniforme, adição/remoção de nós |
TestDataGenerator |
Envio de eventos para o nó DHT correcto via gRPC |
FetchEventsHandler |
Recepção de eventos (storeEvent) e queries (fetchEvents) |
ASTService |
Routing de queries para o nó responsável via DHT |
ASTClient |
Envio de queries para nó aleatório e recepção em stream |
AggregationEngine |
Cálculo de COUNT, SUM, MAX, MIN, SUM_OF_PRODUCTS |
AggregationCache |
Cache hit na segunda query idêntica |
SARouter |
Receção de pedidos ZMQ do SS, despacho async, envio de respostas |
Address already in use ao arrancar o AST
# Windows
netstat -ano | findstr :50051
taskkill /PID <PID> /FResultados sempre 0.0
Os dados de teste podem estar desactualizados. Apaga os CSVs e volta a correr o TestDataGenerator:
# apaga os dados antigos
rm -rf data/
# volta a gerar (com os nós AST em execução)
mvn exec:java -Dexec.mainClass="AST.util.TestDataGenerator"Loop infinito de fetchEvents nos logs do AST
Verifica se o QueryConverter.fromProto está a converter strings vazias para null nos campos fieldK2 e fieldK3.
Cannot resolve symbol após mudanças no .proto
mvn compile| Operação | Descrição | Campos necessários |
|---|---|---|
COUNT |
Conta o número de eventos | — |
SUM |
Soma os valores de um campo | fieldK2 |
MAX |
Valor máximo de um campo | fieldK2 |
MIN |
Valor mínimo de um campo | fieldK2 |
SUM_OF_PRODUCTS |
Soma dos produtos de dois campos | fieldK2, fieldK3 |
Validar o caminho de ingestão de eventos desde o Servidor de Sessão (SS, Erlang) até ao Serviço de Armazenamento de Séries Temporais (AST, Java), passando por ZeroMQ (PUSH/PULL) e serialização Protobuf.
Dispositivo ──TCP──▶ SS (Erlang)
│
│ event:encode_msg (Protobuf)
│ chumak PUSH
▼
AST (Java)
│ JeroMQ PULL
│ EventProto.Event.parseFrom
▼
DHT (consistent hashing)
O SS faz connect do socket PUSH ao único nó AST que conhece (porto 5550 por
defeito). O AST faz bind do PULL nesse porto. A DHT é completamente
transparente ao SS — é o nó AST de contacto que reencaminha internamente para o
nó responsável.
- Ficheiros de teste na pasta
src/do projeto Erlang:ast_link_test.erl— teste isolado (só transporte + serialização)test_ss_ast.erl— teste ponta-a-ponta com o SS real
dispositivos.datna raiz do projeto Erlang, com pelo menos:{<<"drone_01">>, <<"senha123">>, <<"drone">>}.- Script
launch_ast.pypara levantar os nós AST
No terminal, a partir da raiz do projeto Java:
python launch_ast.py 3Esperar que o launcher confirme que todos os nós estão prontos:
[AST-1] Porto 5550 pronto (1.2s)
[AST-2] Porto 5551 pronto (1.3s)
[AST-3] Porto 5552 pronto (1.5s)
[launcher] Todos os 3 nós AST prontos.
Se algum nó não arrancar, o launcher mostra as últimas linhas do log
(ast1.log, ast2.log, ast3.log) para diagnóstico.
Este teste não precisa do SS a correr. Liga um socket PUSH diretamente ao
AST e empurra eventos codificados com o módulo event (gpb). Serve para
confirmar, de forma isolada, que o transporte ZeroMQ funciona e que o Protobuf
desserializa corretamente entre Erlang e Java.
Numa rebar3 shell do projeto Erlang:
1> ast_link_test:fire(3).Resultado esperado no terminal Erlang:
[test] enviado evento #1 (66 bytes) -> 127.0.0.1:5550
[test] enviado evento #2 (66 bytes) -> 127.0.0.1:5550
[test] enviado evento #3 (66 bytes) -> 127.0.0.1:5550
[test] terminado.
Resultado esperado em ast1.log:
[AST-ZMQ] Received Bytes: 66
INFO: Event successfully forwarded to responsible DHT node
[AST-ZMQ] Event stored: drone-1
[AST-ZMQ] Received Bytes: 66
INFO: Event successfully forwarded to responsible DHT node
[AST-ZMQ] Event stored: drone-2
[AST-ZMQ] Received Bytes: 66
INFO: Event successfully forwarded to responsible DHT node
[AST-ZMQ] Event stored: drone-3
| Sintoma | Causa provável | Solução |
|---|---|---|
zmq connection failed no Erlang |
AST não está a correr ou porto errado | Verificar que launch_ast.py mostrou "Porto 5550 pronto" |
Nada aparece no ast1.log |
Slow joiner (mensagem perdida antes do handshake) | Repetir fire(3) — o sleep de 300ms no código já mitiga isto |
FALHA a desserializar no log do AST |
.proto divergente entre Erlang e Java |
Regenerar ambos os lados a partir do mesmo event.proto |
Num segundo terminal, a partir da raiz do projeto Erlang:
rebar3 shell1> session_server:start_server(5555, none, 6001, 5550).Os parâmetros são: porta REP/REQ do SS, seed (none = nó primário), porta do SA, porta ZMQ do AST. Confirmar que aparece:
[Porto 5555] Servidor de Sessao a escuta (SA Router -> 6001)...
E que não há mensagens zmq connection failed (o AST já está de pé desde o
Passo 1, por isso o PUSH liga de imediato).
Este teste percorre o caminho completo: o cliente liga-se ao SS por REQ/REP, faz
login de um dispositivo registado no dispositivos.dat, envia eventos que
passam pelo device_manager (autenticação, construção do mapa Protobuf, encode,
PUSH), e confirma que chegam ao AST.
Num terceiro terminal:
rebar3 shell1> test_ss_ast:run().Resultado esperado no terminal do teste:
========== Teste E2E: SS -> AST ==========
[e2e] Login de drone_01 na porta 5555...
[e2e] Login OK
[e2e] A enviar evento #1 ...
[e2e] Evento #1 -> OK
[e2e] A enviar evento #2 ...
[e2e] Evento #2 -> OK
[e2e] A enviar evento #3 ...
[e2e] Evento #3 -> OK
[e2e] Logout de drone_01...
========== Teste E2E concluído ==========
Resultado esperado no terminal do SS:
[DM] Enviando evento para AST via ZMQ, 67 bytes
[DM] Enviando evento para AST via ZMQ, 67 bytes
[DM] Enviando evento para AST via ZMQ, 67 bytes
Resultado esperado em ast1.log:
[AST-ZMQ] Event stored: drone_01 (×3)
| Sintoma | Causa provável | Solução |
|---|---|---|
Login FALHOU: auth_failed |
Credenciais não batem com dispositivos.dat |
Verificar que o ficheiro tem {<<"drone_01">>, <<"senha123">>, <<"drone">>}. e não está truncado |
not_logged_in no record_event |
Login não foi feito ou expirou | Confirmar que o login devolveu ok antes de enviar eventos |
| Eventos OK no SS mas nada no AST | PUSH ligado a porto errado | Confirmar que o 4º argumento de start_server é 5550 (ou o porto ZMQ real do AST) |
Após o Passo 4, na mesma shell do cliente, enviar um evento intencionalmente inválido (valor inteiro onde o Protobuf espera string) seguido de um evento válido:
%% Evento inválido — valor inteiro em vez de binário
BadEvent = #{<<"type">> => <<"alarme">>, <<"bateria">> => 42},
session_server:client_request(5555, {record_event, <<"drone_01">>, BadEvent}).Resultado esperado: {error, {bad_event, badarg}}
%% Evento válido logo a seguir — confirma que o SS sobreviveu
GoodEvent = #{<<"type">> => <<"alarme">>, <<"bateria">> => <<"42">>},
session_server:client_request(5555, {record_event, <<"drone_01">>, GoodEvent}).Resultado esperado: {ok, event_recorded}
No terminal do SS aparece a linha de erro para o primeiro evento:
[DM] ERRO ao encodar/enviar evento de <<"drone_01">>: error:badarg
E a linha normal para o segundo:
[DM] Enviando evento para AST via ZMQ, 67 bytes
O ponto crítico: o SS continua vivo após o evento malformado. Sem o
try/catch no device_manager, o processo device_manager morria e todos os
clientes subsequentes ficariam sem resposta.
| Ficheiro | Linguagem | Objetivo | Precisa do SS? |
|---|---|---|---|
ast_link_test.erl |
Erlang | Teste isolado: transporte ZMQ + serialização Protobuf | Não |
test_ss_ast.erl |
Erlang | Teste ponta-a-ponta: login + envio de eventos pelo SS real | Sim |
launch_ast.py |
Python | Levantar N nós AST com health check e logs separados | — |
-
Um só
event.proto. Tanto o Erlang (gpb) como o Java (protoc) devem ser gerados a partir do mesmo ficheiro.proto. Se divergirem, o teste isolado (Passo 2) deteta-o como falha de desserialização. -
Valores sempre como binários. No
map<string, string>do Protobuf, tanto as chaves como os valores dos campos do evento têm de ser binários (<<"42">>, não42). Valores inteiros causambadargno encode — otry/catchadicionado no Passo 5 protege contra isto. -
Slow joiner. Um socket PUSH que envia imediatamente após o
connectpode perder a primeira mensagem antes do handshake ZMTP completar. Os módulos de teste incluem umtimer:sleep(300)após o connect para mitigar isto. Em produção, o SS envia eventos apenas após login de dispositivos, o que dá tempo suficiente para a ligação assentar. -
PUSH/PULL descarta em silêncio. Se o AST não estiver acessível e o HWM (high water mark) do PUSH encher, as mensagens são descartadas sem erro. Este comportamento é tema de um teste de robustez mais avançado, fora do âmbito deste documento.