/*
Precisamos transferir dados de uma fonte para um consumidor.
A fonte fornece os dados em pequenos lotes (~dez registos), enquanto que o consumidor funciona melhor com lotes maiores (~mil registos).
Um exemplo real é a transferência de dados de filas tipo Kafka para uma base de dados Clickhouse.
Fonte:
- Praticamente infinita.
- A fonte nunca devolve mais de MaxItems registos numa única chamada a Next.
- Dentro de uma "sessão" (uma chamada à função Pipe), a fonte devolve novos dados em cada chamada a Next.
- No entanto, após reiniciar, a fonte começa a partir da posição "confirmada" anterior, definida por cookie.
Portanto, *cada* valor de cookie devolvido por Next, após guardar os dados no receptor,
deve ser confirmado com uma chamada a Commit, na mesma ordem em que foram devolvidos por Next.
Receptor:
- Não pode processar mais de MaxItems de cada vez.
É necessário implementar a função func Pipe(p Producer, c Consumer) error
que lê dados da fonte, os agrupa num buffer de tamanho não superior a MaxItems e os guarda no receptor,
depois de o fazer, confirma o progresso na fonte.
*/
const MaxItems = 9999
type Producer interface {
// Next devolve:
// - um lote de itens para processar
// - cookie a ser confirmado quando o processamento estiver concluído
// - erro
Next() (items []any, cookie int, err error)
// Commit é usado para marcar o lote de dados como processado
Commit(cookie int) error
}
type Consumer interface {
Process(items []any) error
}
func Pipe(p Producer, c Consumer) error {
// TODO
}