Sobes.tech
Middle+

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