This is an automated email from the ASF dual-hosted git repository.
rong pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/master by this push:
new 32e3d066813 Pipe: rename PipeExtractor -> PipeSource, PipeConnector ->
PipeSink in pipe-api, antlr (#11427)
32e3d066813 is described below
commit 32e3d066813b73a9ee42084513f826472be77834
Author: Steve Yurong Su <[email protected]>
AuthorDate: Tue Oct 31 09:47:23 2023 +0800
Pipe: rename PipeExtractor -> PipeSource, PipeConnector -> PipeSink in
pipe-api, antlr (#11427)
* pipe-api: PipeExtractor -> PipeSource, PipeConnector -> PipeSink
* SHOW ((PIPE pipeName=identifier) | PIPES (WHERE (CONNECTOR | SINK) USED
BY pipeName=identifier)?)
* Update ColumnHeaderConstant.java: PipeExtractor -> PipeSource,
PipeConnector -> PipeSink
---
.../org/apache/iotdb/pipe/api/PipeConnector.java | 3 +-
.../org/apache/iotdb/pipe/api/PipeExtractor.java | 3 +-
.../pipe/api/{PipeConnector.java => PipeSink.java} | 65 +++++++++++-----------
.../api/{PipeExtractor.java => PipeSource.java} | 48 ++++++++--------
.../PipeSinkRuntimeConfiguration.java | 22 ++++++++
.../PipeSourceRuntimeConfiguration.java | 22 ++++++++
.../pipe/api/exception/PipeSinkException.java | 31 +++++++++++
.../org/apache/iotdb/db/qp/sql/IoTDBSqlParser.g4 | 2 +-
.../common/header/ColumnHeaderConstant.java | 4 +-
9 files changed, 137 insertions(+), 63 deletions(-)
diff --git
a/iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/PipeConnector.java
b/iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/PipeConnector.java
index 3ceb6f73f56..ef79ba388cb 100644
---
a/iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/PipeConnector.java
+++
b/iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/PipeConnector.java
@@ -28,7 +28,7 @@ import
org.apache.iotdb.pipe.api.event.dml.insertion.TsFileInsertionEvent;
import org.apache.iotdb.pipe.api.exception.PipeConnectionException;
/**
- * PipeConnector
+ * PipeConnector (Deprecated since v1.3.0, renamed to PipeSink)
*
* <p>PipeConnector is responsible for sending events to sinks.
*
@@ -63,6 +63,7 @@ import
org.apache.iotdb.pipe.api.exception.PipeConnectionException;
* will be called to create a new connection with the sink when the method
{@link
* PipeConnector#heartbeat()} throws exceptions.
*/
+@Deprecated // since v1.3.0, renamed to PipeSink
public interface PipeConnector extends PipePlugin {
/**
diff --git
a/iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/PipeExtractor.java
b/iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/PipeExtractor.java
index a99d3afed9e..87e11c356ee 100644
---
a/iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/PipeExtractor.java
+++
b/iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/PipeExtractor.java
@@ -25,7 +25,7 @@ import
org.apache.iotdb.pipe.api.customizer.parameter.PipeParameters;
import org.apache.iotdb.pipe.api.event.Event;
/**
- * PipeExtractor
+ * PipeExtractor (Deprecated since v1.3.0, renamed to PipeSource)
*
* <p>PipeExtractor is responsible for capturing events from sources.
*
@@ -48,6 +48,7 @@ import org.apache.iotdb.pipe.api.event.Event;
* cancelled (the `DROP PIPE` command is executed).
* </ul>
*/
+@Deprecated // since v1.3.0, renamed to PipeSource
public interface PipeExtractor extends PipePlugin {
/**
diff --git
a/iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/PipeConnector.java
b/iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/PipeSink.java
similarity index 66%
copy from
iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/PipeConnector.java
copy to iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/PipeSink.java
index 3ceb6f73f56..d2a5262ec94 100644
---
a/iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/PipeConnector.java
+++ b/iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/PipeSink.java
@@ -19,7 +19,7 @@
package org.apache.iotdb.pipe.api;
-import
org.apache.iotdb.pipe.api.customizer.configuration.PipeConnectorRuntimeConfiguration;
+import
org.apache.iotdb.pipe.api.customizer.configuration.PipeSinkRuntimeConfiguration;
import org.apache.iotdb.pipe.api.customizer.parameter.PipeParameterValidator;
import org.apache.iotdb.pipe.api.customizer.parameter.PipeParameters;
import org.apache.iotdb.pipe.api.event.Event;
@@ -28,46 +28,44 @@ import
org.apache.iotdb.pipe.api.event.dml.insertion.TsFileInsertionEvent;
import org.apache.iotdb.pipe.api.exception.PipeConnectionException;
/**
- * PipeConnector
+ * PipeSink
*
- * <p>PipeConnector is responsible for sending events to sinks.
+ * <p>PipeSink is responsible for sending events to sinks.
*
- * <p>Various network protocols can be supported by implementing different
PipeConnector classes.
+ * <p>Various network protocols can be supported by implementing different
PipeSink classes.
*
- * <p>The lifecycle of a PipeConnector is as follows:
+ * <p>The lifecycle of a PipeSink is as follows:
*
* <ul>
* <li>When a collaboration task is created, the KV pairs of `WITH
CONNECTOR` clause in SQL are
- * parsed and the validation method {@link
PipeConnector#validate(PipeParameterValidator)}
- * will be called to validate the parameters.
- * <li>Before the collaboration task starts, the method {@link
- * PipeConnector#customize(PipeParameters,
PipeConnectorRuntimeConfiguration)} will be called
- * to config the runtime behavior of the PipeConnector and the method
{@link
- * PipeConnector#handshake()} will be called to create a connection with
sink.
+ * parsed and the validation method {@link
PipeSink#validate(PipeParameterValidator)} will be
+ * called to validate the parameters.
+ * <li>Before the collaboration task starts, the method {@link
PipeSink#customize(PipeParameters,
+ * PipeSinkRuntimeConfiguration)} will be called to config the runtime
behavior of the
+ * PipeSink and the method {@link PipeSink#handshake()} will be called
to create a connection
+ * with sink.
* <li>While the collaboration task is in progress:
* <ul>
* <li>PipeExtractor captures the events and wraps them into three
types of Event instances.
- * <li>PipeProcessor processes the event and then passes them to the
PipeConnector.
- * <li>PipeConnector serializes the events into binaries and send them
to sinks. The
- * following 3 methods will be called: {@link
- * PipeConnector#transfer(TabletInsertionEvent)}, {@link
- * PipeConnector#transfer(TsFileInsertionEvent)} and {@link
- * PipeConnector#transfer(Event)}.
+ * <li>PipeProcessor processes the event and then passes them to the
PipeSink.
+ * <li>PipeSink serializes the events into binaries and send them to
sinks. The following 3
+ * methods will be called: {@link
PipeSink#transfer(TabletInsertionEvent)}, {@link
+ * PipeSink#transfer(TsFileInsertionEvent)} and {@link
PipeSink#transfer(Event)}.
* </ul>
* <li>When the collaboration task is cancelled (the `DROP PIPE` command is
executed), the {@link
- * PipeConnector#close() } method will be called.
+ * PipeSink#close() } method will be called.
* </ul>
*
- * <p>In addition, the method {@link PipeConnector#heartbeat()} will be called
periodically to check
- * whether the connection with sink is still alive. The method {@link
PipeConnector#handshake()}
- * will be called to create a new connection with the sink when the method
{@link
- * PipeConnector#heartbeat()} throws exceptions.
+ * <p>In addition, the method {@link PipeSink#heartbeat()} will be called
periodically to check
+ * whether the connection with sink is still alive. The method {@link
PipeSink#handshake()} will be
+ * called to create a new connection with the sink when the method {@link
PipeSink#heartbeat()}
+ * throws exceptions.
*/
-public interface PipeConnector extends PipePlugin {
+public interface PipeSink extends PipeConnector {
/**
* This method is mainly used to validate {@link PipeParameters} and it is
executed before {@link
- * PipeConnector#customize(PipeParameters,
PipeConnectorRuntimeConfiguration)} is called.
+ * PipeSink#customize(PipeParameters, PipeSinkRuntimeConfiguration)} is
called.
*
* @param validator the validator used to validate {@link PipeParameters}
* @throws Exception if any parameter is not valid
@@ -75,29 +73,28 @@ public interface PipeConnector extends PipePlugin {
void validate(PipeParameterValidator validator) throws Exception;
/**
- * This method is mainly used to customize PipeConnector. In this method,
the user can do the
- * following things:
+ * This method is mainly used to customize PipeSink. In this method, the
user can do the following
+ * things:
*
* <ul>
* <li>Use PipeParameters to parse key-value pair attributes entered by
the user.
- * <li>Set the running configurations in PipeConnectorRuntimeConfiguration.
+ * <li>Set the running configurations in PipeSinkRuntimeConfiguration.
* </ul>
*
- * <p>This method is called after the method {@link
- * PipeConnector#validate(PipeParameterValidator)} is called and before the
method {@link
- * PipeConnector#handshake()} is called.
+ * <p>This method is called after the method {@link
PipeSink#validate(PipeParameterValidator)} is
+ * called and before the method {@link PipeSink#handshake()} is called.
*
* @param parameters used to parse the input parameters entered by the user
- * @param configuration used to set the required properties of the running
PipeConnector
+ * @param configuration used to set the required properties of the running
PipeSink
* @throws Exception the user can throw errors if necessary
*/
- void customize(PipeParameters parameters, PipeConnectorRuntimeConfiguration
configuration)
+ void customize(PipeParameters parameters, PipeSinkRuntimeConfiguration
configuration)
throws Exception;
/**
* This method is used to create a connection with sink. This method will be
called after the
- * method {@link PipeConnector#customize(PipeParameters,
PipeConnectorRuntimeConfiguration)} is
- * called or will be called when the method {@link
PipeConnector#heartbeat()} throws exceptions.
+ * method {@link PipeSink#customize(PipeParameters,
PipeSinkRuntimeConfiguration)} is called or
+ * will be called when the method {@link PipeSink#heartbeat()} throws
exceptions.
*
* @throws Exception if the connection is failed to be created
*/
diff --git
a/iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/PipeExtractor.java
b/iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/PipeSource.java
similarity index 61%
copy from
iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/PipeExtractor.java
copy to
iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/PipeSource.java
index a99d3afed9e..e68724055d6 100644
---
a/iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/PipeExtractor.java
+++ b/iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/PipeSource.java
@@ -19,40 +19,40 @@
package org.apache.iotdb.pipe.api;
-import
org.apache.iotdb.pipe.api.customizer.configuration.PipeExtractorRuntimeConfiguration;
+import
org.apache.iotdb.pipe.api.customizer.configuration.PipeSourceRuntimeConfiguration;
import org.apache.iotdb.pipe.api.customizer.parameter.PipeParameterValidator;
import org.apache.iotdb.pipe.api.customizer.parameter.PipeParameters;
import org.apache.iotdb.pipe.api.event.Event;
/**
- * PipeExtractor
+ * PipeSource
*
- * <p>PipeExtractor is responsible for capturing events from sources.
+ * <p>PipeSource is responsible for capturing events from sources.
*
- * <p>Various data sources can be supported by implementing different
PipeExtractor classes.
+ * <p>Various data sources can be supported by implementing different
PipeSource classes.
*
- * <p>The lifecycle of a PipeExtractor is as follows:
+ * <p>The lifecycle of a PipeSource is as follows:
*
* <ul>
* <li>When a collaboration task is created, the KV pairs of `WITH
EXTRACTOR` clause in SQL are
- * parsed and the validation method {@link
PipeExtractor#validate(PipeParameterValidator)}
- * will be called to validate the parameters.
+ * parsed and the validation method {@link
PipeSource#validate(PipeParameterValidator)} will
+ * be called to validate the parameters.
* <li>Before the collaboration task starts, the method {@link
- * PipeExtractor#customize(PipeParameters,
PipeExtractorRuntimeConfiguration)} will be called
- * to config the runtime behavior of the PipeExtractor.
- * <li>Then the method {@link PipeExtractor#start()} will be called to start
the PipeExtractor.
- * <li>While the collaboration task is in progress, the method {@link
PipeExtractor#supply()} will
- * be called to capture events from sources and then the events will be
passed to the
+ * PipeSource#customize(PipeParameters, PipeSourceRuntimeConfiguration)}
will be called to
+ * config the runtime behavior of the PipeSource.
+ * <li>Then the method {@link PipeSource#start()} will be called to start
the PipeSource.
+ * <li>While the collaboration task is in progress, the method {@link
PipeSource#supply()} will be
+ * called to capture events from sources and then the events will be
passed to the
* PipeProcessor.
- * <li>The method {@link PipeExtractor#close()} will be called when the
collaboration task is
+ * <li>The method {@link PipeSource#close()} will be called when the
collaboration task is
* cancelled (the `DROP PIPE` command is executed).
* </ul>
*/
-public interface PipeExtractor extends PipePlugin {
+public interface PipeSource extends PipeExtractor {
/**
* This method is mainly used to validate {@link PipeParameters} and it is
executed before {@link
- * PipeExtractor#customize(PipeParameters,
PipeExtractorRuntimeConfiguration)} is called.
+ * PipeSource#customize(PipeParameters, PipeSourceRuntimeConfiguration)} is
called.
*
* @param validator the validator used to validate {@link PipeParameters}
* @throws Exception if any parameter is not valid
@@ -60,28 +60,28 @@ public interface PipeExtractor extends PipePlugin {
void validate(PipeParameterValidator validator) throws Exception;
/**
- * This method is mainly used to customize PipeExtractor. In this method,
the user can do the
+ * This method is mainly used to customize PipeSource. In this method, the
user can do the
* following things:
*
* <ul>
* <li>Use PipeParameters to parse key-value pair attributes entered by
the user.
- * <li>Set the running configurations in PipeExtractorRuntimeConfiguration.
+ * <li>Set the running configurations in PipeSourceRuntimeConfiguration.
* </ul>
*
- * <p>This method is called after the method {@link
- * PipeExtractor#validate(PipeParameterValidator)} is called.
+ * <p>This method is called after the method {@link
PipeSource#validate(PipeParameterValidator)}
+ * is called.
*
* @param parameters used to parse the input parameters entered by the user
- * @param configuration used to set the required properties of the running
PipeExtractor
+ * @param configuration used to set the required properties of the running
PipeSource
* @throws Exception the user can throw errors if necessary
*/
- void customize(PipeParameters parameters, PipeExtractorRuntimeConfiguration
configuration)
+ void customize(PipeParameters parameters, PipeSourceRuntimeConfiguration
configuration)
throws Exception;
/**
* Start the extractor. After this method is called, events should be ready
to be supplied by
- * {@link PipeExtractor#supply()}. This method is called after {@link
- * PipeExtractor#customize(PipeParameters,
PipeExtractorRuntimeConfiguration)} is called.
+ * {@link PipeSource#supply()}. This method is called after {@link
+ * PipeSource#customize(PipeParameters, PipeSourceRuntimeConfiguration)} is
called.
*
* @throws Exception the user can throw errors if necessary
*/
@@ -89,7 +89,7 @@ public interface PipeExtractor extends PipePlugin {
/**
* Supply single event from the extractor and the caller will send the event
to the processor.
- * This method is called after {@link PipeExtractor#start()} is called.
+ * This method is called after {@link PipeSource#start()} is called.
*
* @return the event to be supplied. the event may be null if the extractor
has no more events at
* the moment, but the extractor is still running for more events.
diff --git
a/iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/customizer/configuration/PipeSinkRuntimeConfiguration.java
b/iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/customizer/configuration/PipeSinkRuntimeConfiguration.java
new file mode 100644
index 00000000000..40ee575db86
--- /dev/null
+++
b/iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/customizer/configuration/PipeSinkRuntimeConfiguration.java
@@ -0,0 +1,22 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.pipe.api.customizer.configuration;
+
+public interface PipeSinkRuntimeConfiguration extends
PipeConnectorRuntimeConfiguration {}
diff --git
a/iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/customizer/configuration/PipeSourceRuntimeConfiguration.java
b/iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/customizer/configuration/PipeSourceRuntimeConfiguration.java
new file mode 100644
index 00000000000..fd880f30b7f
--- /dev/null
+++
b/iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/customizer/configuration/PipeSourceRuntimeConfiguration.java
@@ -0,0 +1,22 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.pipe.api.customizer.configuration;
+
+public interface PipeSourceRuntimeConfiguration extends
PipeExtractorRuntimeConfiguration {}
diff --git
a/iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/exception/PipeSinkException.java
b/iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/exception/PipeSinkException.java
new file mode 100644
index 00000000000..d947b1cb213
--- /dev/null
+++
b/iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/exception/PipeSinkException.java
@@ -0,0 +1,31 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.pipe.api.exception;
+
+public class PipeSinkException extends PipeConnectionException {
+
+ public PipeSinkException(String message) {
+ super(message);
+ }
+
+ public PipeSinkException(String message, Throwable cause) {
+ super(message, cause);
+ }
+}
diff --git
a/iotdb-core/antlr/src/main/antlr4/org/apache/iotdb/db/qp/sql/IoTDBSqlParser.g4
b/iotdb-core/antlr/src/main/antlr4/org/apache/iotdb/db/qp/sql/IoTDBSqlParser.g4
index 9524a64c7f3..b2224e51eb7 100644
---
a/iotdb-core/antlr/src/main/antlr4/org/apache/iotdb/db/qp/sql/IoTDBSqlParser.g4
+++
b/iotdb-core/antlr/src/main/antlr4/org/apache/iotdb/db/qp/sql/IoTDBSqlParser.g4
@@ -564,7 +564,7 @@ stopPipe
;
showPipes
- : SHOW ((PIPE pipeName=identifier) | PIPES (WHERE CONNECTOR USED BY
pipeName=identifier)?)
+ : SHOW ((PIPE pipeName=identifier) | PIPES (WHERE (CONNECTOR | SINK) USED
BY pipeName=identifier)?)
;
// Pipe Plugin
=========================================================================================
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/common/header/ColumnHeaderConstant.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/common/header/ColumnHeaderConstant.java
index 296a5bc8c1f..4414aace944 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/common/header/ColumnHeaderConstant.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/common/header/ColumnHeaderConstant.java
@@ -150,9 +150,9 @@ public class ColumnHeaderConstant {
// column names for show pipe
public static final String ID = "ID";
public static final String CREATION_TIME = "CreationTime";
- public static final String PIPE_EXTRACTOR = "PipeExtractor";
+ public static final String PIPE_EXTRACTOR = "PipeSource";
public static final String PIPE_PROCESSOR = "PipeProcessor";
- public static final String PIPE_CONNECTOR = "PipeConnector";
+ public static final String PIPE_CONNECTOR = "PipeSink";
public static final String EXCEPTION_MESSAGE = "ExceptionMessage";
// column names for select into