Sobes.tech
Middle+

/* Necesitamos transferir datos de una fuente a un consumidor. La fuente entrega los datos en pequeños lotes (~diez registros), mientras que el consumidor funciona mejor con lotes grandes (~mil registros). Un ejemplo real es la transferencia de datos desde colas tipo Kafka a una base de datos Clickhouse. Fuente: - Prácticamente infinita. - La fuente nunca devuelve más de MaxItems registros en una sola llamada a Next. - Dentro de una "sesión" (una llamada a la función Pipe), la fuente devuelve datos nuevos en cada llamada a Next. - Sin embargo, tras reiniciar, la fuente comienza desde la posición "confirmada" anterior, definida por cookie. Por lo tanto, *cada* valor de cookie que devuelve Next, después de guardar los datos en el receptor, debe ser confirmado con una llamada a Commit, en el mismo orden en que fueron devueltos por Next. Receptor: - No puede procesar más de MaxItems a la vez. Se requiere implementar la función func Pipe(p Producer, c Consumer) error que lee datos de la fuente, los agrupa en un buffer de tamaño no mayor a MaxItems y los guarda en el receptor, después de lo cual confirma el progreso en la fuente. */ const MaxItems = 9999 type Producer interface { // Next devuelve: // - un lote de items para procesar // - cookie para confirmar cuando se complete el procesamiento // - error Next() (items []any, cookie int, err error) // Commit se usa para marcar el lote de datos como procesado Commit(cookie int) error } type Consumer interface { Process(items []any) error } func Pipe(p Producer, c Consumer) error { // TODO }