You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
Copy file name to clipboardExpand all lines: flink-connector-elasticsearch-base/src/main/java/org/apache/flink/streaming/connectors/elasticsearch/table/ElasticsearchConfiguration.java
+4
Original file line number
Diff line number
Diff line change
@@ -150,6 +150,10 @@ public Optional<String> getPathPrefix() {
Copy file name to clipboardExpand all lines: flink-connector-elasticsearch-base/src/main/java/org/apache/flink/streaming/connectors/elasticsearch/table/ElasticsearchConnectorOptions.java
+6
Original file line number
Diff line number
Diff line change
@@ -152,6 +152,12 @@ public class ElasticsearchConnectorOptions {
152
152
"The format must produce a valid JSON document. "
153
153
+ "Please refer to the documentation on formats for more details.");
Copy file name to clipboardExpand all lines: flink-connector-elasticsearch-base/src/main/java/org/apache/flink/streaming/connectors/elasticsearch/table/KeyExtractor.java
Copy file name to clipboardExpand all lines: flink-connector-elasticsearch-base/src/main/java/org/apache/flink/streaming/connectors/elasticsearch/table/RequestFactory.java
Copy file name to clipboardExpand all lines: flink-connector-elasticsearch-base/src/main/java/org/apache/flink/streaming/connectors/elasticsearch/table/RowElasticsearchSinkFunction.java
+11-6
Original file line number
Diff line number
Diff line change
@@ -51,20 +51,23 @@ class RowElasticsearchSinkFunction implements ElasticsearchSinkFunction<RowData>
51
51
privatefinalXContentTypecontentType;
52
52
privatefinalRequestFactoryrequestFactory;
53
53
privatefinalFunction<RowData, String> createKey;
54
+
privatefinalFunction<RowData, String> routingKey;
54
55
55
56
publicRowElasticsearchSinkFunction(
56
57
IndexGeneratorindexGenerator,
57
58
@NullableStringdocType, // this is deprecated in es 7+
Copy file name to clipboardExpand all lines: flink-connector-elasticsearch-base/src/test/java/org/apache/flink/streaming/connectors/elasticsearch/table/KeyExtractorTest.java
Copy file name to clipboardExpand all lines: flink-connector-elasticsearch6/src/main/java/org/apache/flink/streaming/connectors/elasticsearch/table/Elasticsearch6DynamicSink.java
+21-7
Original file line number
Diff line number
Diff line change
@@ -148,7 +148,8 @@ public SinkFunctionProvider getSinkRuntimeProvider(Context context) {
Copy file name to clipboardExpand all lines: flink-connector-elasticsearch7/src/main/java/org/apache/flink/streaming/connectors/elasticsearch/table/Elasticsearch7DynamicSink.java
+21-7
Original file line number
Diff line number
Diff line change
@@ -143,7 +143,8 @@ public SinkFunctionProvider getSinkRuntimeProvider(Context context) {
Copy file name to clipboardExpand all lines: flink-connector-elasticsearch7/src/main/java/org/apache/flink/streaming/connectors/elasticsearch/table/Elasticsearch7DynamicTableFactory.java
0 commit comments