This is an automated email from the ASF dual-hosted git repository.
guoweijie pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/flink-connector-pulsar.git
The following commit(s) were added to refs/heads/main by this push:
new fb98096 [FLINK-31048][Connector/Pulsar] Fix incorrect java doc in
PulsarSource and PulsarSourceBuilder.
fb98096 is described below
commit fb98096d9a67d26b1965b8ef181125317a59a5c4
Author: Weijie Guo <[email protected]>
AuthorDate: Tue Feb 14 01:18:06 2023 +0800
[FLINK-31048][Connector/Pulsar] Fix incorrect java doc in PulsarSource and
PulsarSourceBuilder.
This closes #28
---
.../java/org/apache/flink/connector/pulsar/source/PulsarSource.java | 2 +-
.../org/apache/flink/connector/pulsar/source/PulsarSourceBuilder.java | 4 ++--
2 files changed, 3 insertions(+), 3 deletions(-)
diff --git
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/PulsarSource.java
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/PulsarSource.java
index c7d272b..c038f32 100644
---
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/PulsarSource.java
+++
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/PulsarSource.java
@@ -57,7 +57,7 @@ import org.apache.pulsar.client.api.PulsarClientException;
* .setServiceUrl(getServiceUrl())
* .setAdminUrl(getAdminUrl())
* .setSubscriptionName("test")
- * .setDeserializationSchema(PulsarDeserializationSchema.flinkSchema(new
SimpleStringSchema()))
+ * .setDeserializationSchema(new SimpleStringSchema())
* .setBounded(StopCursor::defaultStopCursor)
* .build();
* }</pre>
diff --git
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/PulsarSourceBuilder.java
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/PulsarSourceBuilder.java
index 80d8c30..1c8cfbe 100644
---
a/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/PulsarSourceBuilder.java
+++
b/flink-connector-pulsar/src/main/java/org/apache/flink/connector/pulsar/source/PulsarSourceBuilder.java
@@ -89,7 +89,7 @@ import static org.apache.flink.util.Preconditions.checkState;
* .setAdminUrl(PULSAR_BROKER_HTTP_URL)
* .setSubscriptionName("flink-source-1")
* .setTopics(Arrays.asList(TOPIC1, TOPIC2))
- * .setDeserializationSchema(PulsarDeserializationSchema.flinkSchema(new
SimpleStringSchema()))
+ * .setDeserializationSchema(new SimpleStringSchema())
* .build();
* }</pre>
*
@@ -118,7 +118,7 @@ import static
org.apache.flink.util.Preconditions.checkState;
* .setAdminUrl(PULSAR_BROKER_HTTP_URL)
* .setSubscriptionName("flink-source-1")
* .setTopics(Arrays.asList(TOPIC1, TOPIC2))
- * .setDeserializationSchema(PulsarDeserializationSchema.flinkSchema(new
SimpleStringSchema()))
+ * .setDeserializationSchema(new SimpleStringSchema())
*
.setUnboundedStopCursor(StopCursor.atEventTime(System.currentTimeMillis()))
* .build();
* }</pre>