/*
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
}