Una suscripción WebSocket es un flujo continuo sin cursor: el servidor envía nuevos eventos pero nunca reproduce lo que te perdiste, y un frame de cierre no indica cuántos mensajes se perdieron. La disciplina correcta es tratar el número de bloque de cada notificación (y el hash de bloque para seguridad ante reorganizaciones) como un punto de control monotónico, persistir el último cursor visto y, al reconectar, ejecutar una consulta de relleno acotada sobre la brecha antes de volver a confiar en el flujo en vivo. Como el relleno y la entrega en vivo pueden solaparse, debes deduplicar por (hash de tx, índice de log) o (número de bloque, índice de log) y almacenar en búfer los eventos en vivo hasta alcanzar la marca de agua alta del relleno. El mismo patrón se generaliza a Ethereum, Solana, Sui y Substrate con distintas primitivas. Este artículo ofrece una implementación ejecutable en TypeScript, un bucle de verificación, una tabla de modos de fallo y las ventajas y desventajas frente al sondeo simple.
Por qué una suscripción reconectada pierde eventos silenciosamente
Una suscripción WebSocket es un flujo continuo sin cursor. Según la especificación PubSub de JSON-RPC de Ethereum, eth_subscribe con un parámetro logs devuelve un ID de suscripción y luego envía notificaciones a medida que se minan nuevos bloques. El servidor no reproduce los eventos que ocurrieron mientras tu socket estaba caído, y el frame de cierre no te da un recuento de los mensajes perdidos. Por lo tanto, una reconexión ingenua produce un flujo con un hueco.
Esto es sustancialmente diferente de una llamada RPC de solicitud/respuesta, donde un fallo es visible como un error. Un fallo de suscripción es invisible: tu manejador simplemente deja de recibir y luego se reanuda en algún bloque posterior. Si estás indexando transferencias, liquidaciones o eventos de gobernanza, un flujo de logs con un hueco es peor que no tener flujo, porque los consumidores posteriores asumen que está completo.
Las causas de la desconexión en sí se tratan en Por qué se desconectan los WebSocket RPC y cómo solucionarlo; este artículo asume que ya reconectas y se centra en el problema más difícil: demostrar que no te perdiste nada y recuperar la brecha de forma determinista.
- Las suscripciones solo envían nuevos eventos; no hay búfer de reproducción en el servidor.
- Un frame de cierre no lleva número de secuencia ni recuento de mensajes perdidos.
- Reconectar sin reconciliación = brecha silenciosa.
- Los consumidores críticos para la corrección deben tratar la suscripción como una pista, no como un libro mayor.
El modelo mental del cursor: el número de bloque como punto de control monotónico
Trata el número de bloque de cada notificación como un punto de control monotónico. Persiste el número de bloque más alto que hayas procesado por completo como lastCursor. Para diseños seguros ante reorganizaciones, persiste también el hash de bloque para poder detectar una reorganización de la cadena que invalide tu cursor, siguiendo las convenciones de identidad de bloque descritas en EIP-1898 y la semántica de número de bloque de EIP-234.
Al reconectar, la regla es: no confíes en el flujo en vivo hasta que hayas rellenado [lastCursor + 1 .. latest]. El relleno es una consulta eth_getLogs acotada sobre ese rango. Solo después de que el relleno se complete y hayas avanzado lastCursor hasta la marca de agua alta del relleno deberías reanudar el consumo de notificaciones en vivo.
Esto convierte un canal push poco fiable en una tubería fiable de pull más push. La suscripción te da baja latencia; el relleno te da completitud. Ninguno por sí solo es suficiente.
- Persiste
lastCursor(número de bloque) y opcionalmentelastHash(hash de bloque). - El rango de relleno es
[lastCursor + 1 .. latest], acotado porlatest. - Avanza el cursor solo después de que el relleno se procese de forma duradera.
- Seguridad ante reorganizaciones: si
lastHashya no coincide con la cadena, retrocede el cursor.
El peligro de ordenación: almacenar en búfer los eventos en vivo hasta que el relleno se complete
Un modo de fallo sutil: después de volver a suscribirte, el flujo en vivo puede empezar a entregar bloques que están por delante de tu marca de agua alta del relleno. Si procesas los eventos en vivo de inmediato, puedes procesar el bloque 100 (en vivo) antes que el bloque 95 (relleno), produciendo un estado desordenado y doble conteo.
La regla es almacenar en búfer los eventos en vivo hasta que el relleno se complete. Recoge las notificaciones entrantes en una cola, ejecuta el relleno, luego fusiona la cola con los resultados del relleno, deduplica y procesa en orden de bloque. Solo después de que la cola se vacíe y el cursor avance cambias a procesamiento en vivo directo.
Esta ventana de almacenamiento en búfer es corta (segundos), pero es la diferencia entre un indexador correcto y uno que ocasionalmente cuenta dos veces una transferencia. La misma disciplina se aplica ya sea en Ethereum o en cualquier otra cadena.
- El flujo en vivo puede empezar por delante de la marca de agua alta del relleno.
- Almacena en búfer los eventos en vivo en una cola durante el relleno.
- Fusiona, deduplica, ordena por número de bloque y luego procesa.
- Cambia a procesamiento en vivo directo solo después de que la cola se vacíe.
Deduplicar por (hash de tx, índice de log) o (número de bloque, índice de log)
El relleno y la entrega en vivo pueden solaparse. Un log que llegó en vivo en el bloque 100 también puede aparecer en tu relleno de eth_getLogs si el rango de relleno se extendió hasta el bloque 100. Sin deduplicación, lo procesas dos veces.
La clave de deduplicación canónica para logs de Ethereum es (transactionHash, logIndex). Para cadenas sin índice de log, usa (blockNumber, eventIndex) o un ID de evento determinista. Mantén un conjunto de deduplicación de corta duración (o una restricción única en tu base de datos) que cubra al menos la ventana de relleno más un margen de seguridad.
La deduplicación no es opcional. Es el mecanismo que hace seguro el solapamiento entre push y pull. Si la omites, tu relleno introduce el mismo doble conteo que pretendía evitar.
- Ethereum: clave de deduplicación
(transactionHash, logIndex). - Genérico:
(blockNumber, eventIndex)o un ID de evento determinista. - Mantén el conjunto de deduplicación al menos tan largo como la ventana de relleno.
- Una restricción única en la base de datos es la deduplicación más robusta.
TypeScript ejecutable: suscribir, detectar reconexión, rellenar, deduplicar, reanudar
El siguiente ejemplo usa ethers v6 y un cursor persistente. Se suscribe a logs, persiste lastBlock, detecta la reconexión, ejecuta un relleno acotado con getLogs, deduplica y solo entonces reanuda el procesamiento en vivo. Reemplaza el endpoint con la URL WebSocket de tu proveedor; el comportamiento específico del proveedor, como los tiempos de espera por inactividad y los rangos máximos de relleno, está documentado / varía según el proveedor.
Observa el array buffer: los eventos en vivo que llegan durante el relleno se encolan, no se procesan. El conjunto seen deduplica por (transactionHash, logIndex). La función backfill está acotada por latest y debe tener en cuenta los límites de tasa.
import { ethers } from "ethers";
const WS_URL = process.env.WS_URL!;
const ADDRESS = process.env.ADDRESS!; // contract to watch
const MAX_RANGE = 2000; // bounded backfill window
let lastBlock = Number(process.env.LAST_BLOCK ?? 0);
let backfilling = false;
const buffer: ethers.Log[] = [];
const seen = new Set<string>();
function key(l: ethers.Log) {
return `${l.transactionHash}:${l.index}`;
}
async function processLog(l: ethers.Log) {
const k = key(l);
if (seen.has(k)) return;
seen.add(k);
// TODO: persist to your store
console.log("processed", k, "block", l.blockNumber);
if (l.blockNumber > lastBlock) lastBlock = l.blockNumber;
}
async function backfill(provider: ethers.Provider) {
backfilling = true;
const latest = await provider.getBlockNumber();
let from = lastBlock + 1;
while (from <= latest) {
const to = Math.min(from + MAX_RANGE - 1, latest);
const logs = await provider.getLogs({ address: ADDRESS, fromBlock: from, toBlock: to });
logs.sort((a, b) => a.blockNumber - b.blockNumber || a.index - b.index);
for (const l of logs) await processLog(l);
from = to + 1;
}
// drain buffered live events
buffer.sort((a, b) => a.blockNumber - b.blockNumber || a.index - b.index);
for (const l of buffer) await processLog(l);
buffer.length = 0;
backfilling = false;
}
async function main() {
const provider = new ethers.WebSocketProvider(WS_URL);
provider.on("error", () => {});
provider.websocket.on("close", async () => {
console.warn("socket closed; reconnecting");
await backfill(provider);
});
provider.on({ address: ADDRESS }, async (l: ethers.Log) => {
if (backfilling) buffer.push(l);
else await processLog(l);
});
await backfill(provider); // initial catch-up
}
main().catch(console.error);Bucle de verificación: mata el socket, inyecta un evento, asegura recuperación exactamente una vez
No puedes confiar en un diseño de recuperación de brechas que no hayas probado. Construye un bucle de verificación que mate deliberadamente el socket, inyecte un evento conocido durante la interrupción y asegure que el evento se recupera exactamente una vez. Esta es la única forma de demostrar que tu lógica de relleno y deduplicación realmente funciona.
El bucle: (1) inicia el suscriptor y registra lastBlock; (2) fuerza el cierre del WebSocket; (3) mientras está desconectado, envía una transacción que emita un evento conocido; (4) permite que la reconexión y el relleno se ejecuten; (5) asegura que el evento aparece exactamente una vez en tu almacén y que lastBlock avanzó más allá de él.
Ejecuta este bucle contra tu propio endpoint y registra los resultados en una tabla. No confíes en las cifras de latencia o fiabilidad publicadas por el proveedor; mide las tuyas.
- Fuerza el cierre del socket programáticamente (p. ej.,
provider.websocket.close()). - Inyecta un evento conocido durante la ventana de interrupción.
- Asegura la recuperación exactamente una vez y el avance del cursor.
- Repite en tiempos de espera por inactividad, despliegues del proveedor y reinicios del balanceador de carga.
Tabla de resultados: mide la recuperación de brechas contra tu propio endpoint
Usa la siguiente tabla para registrar tus propias mediciones. Rellénala con resultados de tu endpoint y tu carga de trabajo. Las cifras específicas del proveedor están documentadas / varían según el proveedor, por lo que tus propias mediciones son la única guía fiable.
Ejecuta cada escenario al menos diez veces y registra el peor caso, no el promedio. La ventana de brecha es lo que importa: si tu rango de relleno supera el rango máximo de eth_getLogs del proveedor, debes dividirlo en fragmentos.
- Escenario | Causa de desconexión | Brecha (bloques) | Tiempo de relleno (ms) | Eventos recuperados | Duplicados | Aprobado/Fallido
- Tiempo de espera por inactividad | Sin tráfico durante N minutos | | | | |
- Despliegue del proveedor | Reinicio del lado del servidor | | | | |
- Reinicio del balanceador de carga | Conexión caída | | | | |
- Cierre forzado | Muerte del lado del cliente | | | | |
- Reorganización | Reorganización de la cadena | | | | |
Generalizar el patrón: Solana, Sui y Substrate
La misma disciplina se aplica en todas las cadenas con distintas primitivas. En Solana, usa un cursor basado en slot: suscríbete a un programa o cuenta, persiste el último slot procesado y al reconectar rellena con getSignaturesForAddress y getTransaction. En Sui, usa un cursor de checkpoint y suix_queryEvents para rellenar la brecha. En Substrate, usa un cursor de número de bloque y system_events en cada bloque.
Las primitivas difieren, pero el modelo mental es idéntico: cursor monotónico, relleno acotado, deduplicación, almacenar en búfer los eventos en vivo hasta que el relleno se complete. Si entiendes el caso de Ethereum, los entiendes todos.
Para la selección de endpoints y la conmutación por error en estas cadenas, consulta la guía de endpoints RPC y Monitorización de nodos RPC, métricas y conmutación por error.
- Solana: cursor de slot + relleno con
getSignaturesForAddress. - Sui: cursor de checkpoint + relleno con
suix_queryEvents. - Substrate: cursor de número de bloque + relleno con
system_events. - Misma disciplina, distintas primitivas.
Acotar el relleno: límites de tasa, retroceso y la ventana de brecha
La ventana de brecha importa. Los tiempos de espera por inactividad, los despliegues del proveedor y los reinicios del balanceador de carga pueden producir brechas que van de segundos a minutos. Tu rango de relleno debe estar acotado y tener en cuenta los límites de tasa, o alcanzarás los límites del proveedor y no podrás recuperarte.
Divide tus llamadas a eth_getLogs en fragmentos (p. ej., 2000 bloques por llamada) y aplica retroceso exponencial en errores de límite de tasa. El retroceso interactúa con tu lógica de reconexión: si reconectas con demasiada agresividad, puedes reconectar a un estado limitado por tasa y fallar de nuevo. Consulta Errores de tiempo de espera RPC: causas y soluciones para patrones de retroceso.
Si tu brecha supera el rango máximo de relleno del proveedor, debes dividirla en fragmentos. Si supera tu ventana de retención, debes recurrir a una resincronización completa. Conoce tus límites antes de necesitarlos.
- Divide las llamadas de relleno para respetar los rangos máximos del proveedor.
- Aplica retroceso exponencial en errores de límite de tasa.
- No reconectes agresivamente a un estado limitado por tasa.
- Conoce tu ventana de retención; recurre a resincronización completa si se supera.
WebSocket con relleno vs sondeo simple: elegir por corrección
El sondeo simple (eth_getLogs con temporizador) es más sencillo e inherentemente libre de brechas si persistes el cursor, pero añade latencia y puede ser más costoso a alta frecuencia. WebSocket con relleno te da baja latencia más completitud, a costa de un código más complejo.
Elige WebSocket con relleno cuando la latencia importe y puedas implementar correctamente la deduplicación y el almacenamiento en búfer. Elige sondeo simple cuando la corrección sea primordial y la latencia tolerable, o cuando la fiabilidad del WebSocket de tu proveedor sea incierta. La comparación se cubre en eth_subscribe logs vs filtros de sondeo.
A menudo lo mejor es un híbrido: WebSocket para pistas de baja latencia, más una reconciliación periódica por sondeo como red de seguridad. Esto captura cualquier brecha que tu lógica de reconexión haya pasado por alto.
- Sondeo: más sencillo, libre de brechas con cursor, mayor latencia.
- WebSocket con relleno: baja latencia, más complejo.
- Híbrido: pistas por WebSocket + reconciliación periódica por sondeo.
- Elige según la tolerancia a la latencia y los requisitos de corrección.
Modos de fallo y solución de problemas
La siguiente tabla enumera modos de fallo comunes y sus soluciones. Úsala cuando tu recuperación de brechas no funcione como se espera.
Si ves duplicados, tu clave de deduplicación es incorrecta o tu conjunto de deduplicación es demasiado corto. Si ves eventos faltantes, tu rango de relleno es incorrecto o tu cursor avanzó prematuramente. Si ves procesamiento desordenado, no estás almacenando en búfer los eventos en vivo durante el relleno.
- Duplicados | Clave de deduplicación incorrecta o conjunto demasiado corto | Usa
(txHash, logIndex), amplía el conjunto. - Eventos faltantes | Rango de relleno incorrecto o cursor avanzado antes de tiempo | Verifica
[lastCursor+1 .. latest], avanza después del relleno. - Desordenado | Eventos en vivo no almacenados en búfer | Almacena en búfer durante el relleno, ordena por bloque.
- Limitado por tasa | Relleno demasiado agresivo | Divide llamadas, retroceso exponencial.
- Corrupción por reorganización | Sin verificación de hash de bloque | Persiste
lastHash, retrocede si no coincide. - Brecha silenciosa | Sin relleno en absoluto | Implementa cursor + relleno.
Limitaciones, ventajas y desventajas, y próximos pasos
Este patrón tiene limitaciones. Asume que tu proveedor admite eth_getLogs sobre el rango de la brecha; algunos proveedores limitan el rango o el número de resultados. Asume que tu cursor es duradero; si tu proceso falla entre el procesamiento y la persistencia, puedes reprocesar u omitir. Asume que no hay reorganizaciones profundas más allá de tu ventana de retención; las reorganizaciones más profundas requieren una resincronización completa.
La contrapartida es la complejidad: estás construyendo un pequeño motor de reconciliación, no solo un suscriptor. Para consumidores críticos para la corrección, esa complejidad está justificada. Para paneles de bajo riesgo, el sondeo simple puede ser suficiente.
Próximos pasos: implementa el bucle de verificación contra tu propio endpoint, rellena la tabla de resultados y revisa Precios de RPC y el Servicio de API para entender las implicaciones de coste. Para una visión más amplia, consulta el centro de aprendizaje de OnFinality.
Condiciones de actualización: revisa este diseño cuando tu proveedor cambie los rangos máximos de relleno, cuando añadas una nueva cadena o cuando observes una reorganización más profunda que tu ventana de retención.
- El rango máximo de relleno del proveedor puede obligar a dividir en fragmentos.
- La durabilidad del cursor requiere persistencia transaccional.
- Las reorganizaciones profundas más allá de la retención requieren resincronización completa.
- Revisa ante cambios del proveedor, nuevas cadenas o reorganizaciones profundas.