public class CloudPubSubSinkTask extends SinkTaskA SinkTask used by a CloudPubSubSinkConnector to write messages to Google Cloud Pub/Sub.
Constructors
CloudPubSubSinkTask()
public CloudPubSubSinkTask()CloudPubSubSinkTask(Publisher publisher)
public CloudPubSubSinkTask(Publisher publisher)| Parameter | |
|---|---|
| Name | Description |
publisher |
com.google.cloud.pubsub.v1.Publisher |
Methods
flush(Map<TopicPartition,OffsetAndMetadata> partitionOffsets)
public void flush(Map<TopicPartition,OffsetAndMetadata> partitionOffsets)| Parameter | |
|---|---|
| Name | Description |
partitionOffsets |
Map<org.apache.kafka.common.TopicPartition,org.apache.kafka.clients.consumer.OffsetAndMetadata> |
org.apache.kafka.connect.sink.SinkTask.flush(java.util.Map<org.apache.kafka.common.TopicPartition,org.apache.kafka.clients.consumer.OffsetAndMetadata>)
put(Collection<SinkRecord> sinkRecords)
public void put(Collection<SinkRecord> sinkRecords)| Parameter | |
|---|---|
| Name | Description |
sinkRecords |
Collection<org.apache.kafka.connect.sink.SinkRecord> |
org.apache.kafka.connect.sink.SinkTask.put(java.util.Collection<org.apache.kafka.connect.sink.SinkRecord>)
start(Map<String,String> props)
public void start(Map<String,String> props)| Parameter | |
|---|---|
| Name | Description |
props |
Map<String,String> |
org.apache.kafka.connect.sink.SinkTask.start(java.util.Map<java.lang.String,java.lang.String>)
stop()
public void stop()org.apache.kafka.connect.sink.SinkTask.stop()
version()
public String version()| Returns | |
|---|---|
| Type | Description |
String |
|