-
Notifications
You must be signed in to change notification settings - Fork 97
/
KafkaTopics.scala
59 lines (50 loc) · 1.64 KB
/
KafkaTopics.scala
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
/*
* Copyright 2018-2022 OVO Energy Limited
*
* SPDX-License-Identifier: Apache-2.0
*/
package fs2.kafka.consumer
import org.apache.kafka.common.PartitionInfo
import scala.concurrent.duration.FiniteDuration
import org.apache.kafka.common.TopicPartition
trait KafkaTopics[F[_]] {
/**
* Returns the partitions for the specified topic.
*
* Timeout is determined by `default.api.timeout.ms`, which
* is set using [[ConsumerSettings#withDefaultApiTimeout]].
*/
def partitionsFor(topic: String): F[List[PartitionInfo]]
/**
* Returns the partitions for the specified topic.
*/
def partitionsFor(topic: String, timeout: FiniteDuration): F[List[PartitionInfo]]
/**
* Returns the first offset for the specified partitions.<br>
* <br>
* Timeout is determined by `default.api.timeout.ms`, which
* is set using [[ConsumerSettings#withDefaultApiTimeout]].
*/
def beginningOffsets(partitions: Set[TopicPartition]): F[Map[TopicPartition, Long]]
/**
* Returns the first offset for the specified partitions.
*/
def beginningOffsets(
partitions: Set[TopicPartition],
timeout: FiniteDuration
): F[Map[TopicPartition, Long]]
/**
* Returns the last offset for the specified partitions.<br>
* <br>
* Timeout is determined by `request.timeout.ms`, which
* is set using [[ConsumerSettings#withRequestTimeout]].
*/
def endOffsets(partitions: Set[TopicPartition]): F[Map[TopicPartition, Long]]
/**
* Returns the last offset for the specified partitions.
*/
def endOffsets(
partitions: Set[TopicPartition],
timeout: FiniteDuration
): F[Map[TopicPartition, Long]]
}