DanielLeens commented on code in PR #12281: URL: https://github.com/apache/seatunnel/pull/12281#discussion_r4003906471
########## seatunnel-connectors-v2/connector-http/connector-http-splunk/src/main/java/org/apache/seatunnel/connectors/seatunnel/splunk/SplunkSourceFactory.java: ########## @@ -0,0 +1,50 @@ +/* + * 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.seatunnel.connectors.seatunnel.splunk; + +import org.apache.seatunnel.api.configuration.util.OptionRule; +import org.apache.seatunnel.api.source.SeaTunnelSource; +import org.apache.seatunnel.api.source.SourceSplit; +import org.apache.seatunnel.api.table.connector.TableSource; +import org.apache.seatunnel.api.table.factory.Factory; +import org.apache.seatunnel.api.table.factory.TableSourceFactoryContext; +import org.apache.seatunnel.connectors.seatunnel.http.source.HttpSourceFactory; +import org.apache.seatunnel.connectors.seatunnel.splunk.config.SplunkSourceOptions; + +import com.google.auto.service.AutoService; + +import java.io.Serializable; + +@AutoService(Factory.class) +public class SplunkSourceFactory extends HttpSourceFactory { + @Override + public String factoryIdentifier() { + return "Splunk"; + } + + @Override + public <T, SplitT extends SourceSplit, StateT extends Serializable> + TableSource<T, SplitT, StateT> createSource(TableSourceFactoryContext context) { + return () -> (SeaTunnelSource<T, SplitT, StateT>) new SplunkSource(context.getOptions()); + } + + @Override + public OptionRule optionRule() { + return getHttpBuilder().required(SplunkSourceOptions.API_KEY).build(); Review Comment: Confirmed against the current head (`90d798ad`). `SplunkSourceFactory.optionRule()` (`SplunkSourceFactory.java:47-49`) only adds `.required(SplunkSourceOptions.API_KEY)` on top of `getHttpBuilder()`, and `getHttpBuilder()` (`HttpSourceFactory.java:54-75`) does not list `KEEP_PARAMS_AS_FORM` either. `SplunkSourceParameter.buildWithConfig()` reads it directly off `pluginConfig.get(SplunkSourceOptions.KEEP_PARAMS_AS_FORM)` (`SplunkSourceParameter.java:49`), and the doc's own example config sets it explicitly. On the real source-parsing path, `ConfigValidationUtils` calls `ConfigValidator.validateUnknownKeys(...)` against `sourceFactory.optionRule()` before the source is even constructed (`ConfigValidationUtils.java:247-250`), so this key is rejected as unrecognized for any job that follows the documented example. Good catch — my last review verified functional behavior end-to-end but didn't drive the documented config through actual config validation. Treating this as a real, current blocker: `optionRule()` needs `.optional(SplunkSourceOptions.KEEP_PARAMS_AS_FORM)`. ########## docs/en/connectors/source/Http-Splunk.md: ########## @@ -0,0 +1,43 @@ +import ChangeLog from '../changelog/connector-http-splunk.md'; + +# Http-Splunk Source Connector + +## Description + +The `Http-Splunk` connector allows batch reading data from Splunk REST API endpoints using an HTTP source pattern. It utilizes token-based authentication and supports synchronous search export endpoints via POST requests. + +## Key Features + +* HTTP-based ingestion from Splunk REST API (`/services/search/v2/jobs/export`) +* Token-based authentication via the `Authorization` header +* Form-urlencoded parameter mapping for search queries and output formats + +## Options + +| Name | Type | Required | Default | Description | +| --- | --- | --- | --- | --- | +| url | String | Yes | - | Splunk REST API search export URL (`/services/search/v2/jobs/export`) | +| api_key | String | Yes | - | Splunk authentication token (e.g., `Splunk <token>`) | +| method | String | No | `POST` | HTTP request method | +| keep_params_as_form | Boolean | No | `true` | Keep parameters as form urlencoded | +| params | Map | Yes | - | Request parameters including `search` and `output_mode` | + +## Example Configuration + +```hocon +source { + Http-Splunk { Review Comment: Confirmed. `SplunkSourceFactory.factoryIdentifier()` returns `"Splunk"` (`SplunkSourceFactory.java:36-38`), and `plugin-mapping.properties:198` registers `seatunnel.source.Splunk = connector-http-splunk` — there is no `Http-Splunk` factory registered anywhere in the tree. The example block at `docs/en/connectors/source/Http-Splunk.md:29` (`Http-Splunk {`) will fail factory discovery exactly as you describe; it needs to be `Splunk {` to match the registered identifier. I verified the registration chain in my last review (factory/plugin-mapping/plugin_config/dist dependency) but never actually parsed the doc's own example config against it. Thanks for catching this — agreed it's a real blocker, the shipped example currently can't run. ########## seatunnel-connectors-v2/connector-http/connector-http-splunk/src/main/java/org/apache/seatunnel/connectors/seatunnel/splunk/config/SplunkSourceParameter.java: ########## @@ -0,0 +1,53 @@ +/* + * 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.seatunnel.connectors.seatunnel.splunk.config; + +import org.apache.seatunnel.api.configuration.ReadonlyConfig; +import org.apache.seatunnel.connectors.seatunnel.http.config.HttpParameter; + +import java.util.HashMap; + +public class SplunkSourceParameter extends HttpParameter { + /** + * Overrides buildWithConfig to accept an explicit apiKey parameter. Splunk's REST API requires + * the API key to be passed specifically as an Authorization header, so this method ensures the + * key is properly extracted and configured. + */ + public void buildWithConfig(ReadonlyConfig pluginConfig, String apiKey) { + super.buildWithConfig(pluginConfig); Review Comment: Confirmed. `HttpSourceOptions.METHOD` defaults to `GET` (`HttpSourceOptions.java:99-103`), and `SplunkSourceParameter.buildWithConfig()` never overrides it — it only sets headers/params/`keepParamsAsForm`/`enableMultilines` (`SplunkSourceParameter.java:31-52`), so an omitted `method` really does resolve to `GET` at runtime. That contradicts the doc's own options table (`method | ... | Default: POST`), and the shipped example only avoids the bug because it sets `method = "POST"` explicitly, masking it. Agreed this is a real functional gap for any user who leaves `method` unset. I'd lean toward the connector hard-defaulting to POST for this endpoint (or validating and rejecting GET) rather than relying on every user to set it explicitly. Good find. ########## seatunnel-connectors-v2/connector-http/connector-http-splunk/src/main/java/org/apache/seatunnel/connectors/seatunnel/splunk/SplunkSource.java: ########## @@ -0,0 +1,56 @@ +/* + * 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.seatunnel.connectors.seatunnel.splunk; + +import org.apache.seatunnel.api.configuration.ReadonlyConfig; +import org.apache.seatunnel.api.table.type.SeaTunnelRow; +import org.apache.seatunnel.connectors.seatunnel.common.source.AbstractSingleSplitReader; +import org.apache.seatunnel.connectors.seatunnel.common.source.SingleSplitReaderContext; +import org.apache.seatunnel.connectors.seatunnel.http.source.HttpSource; +import org.apache.seatunnel.connectors.seatunnel.http.source.HttpSourceReader; +import org.apache.seatunnel.connectors.seatunnel.splunk.config.SplunkSourceOptions; +import org.apache.seatunnel.connectors.seatunnel.splunk.config.SplunkSourceParameter; + +import lombok.extern.slf4j.Slf4j; + +@Slf4j +public class SplunkSource extends HttpSource { + private final SplunkSourceParameter splunkSourceParameter = new SplunkSourceParameter(); + + public SplunkSource(ReadonlyConfig pluginConfig) { + super(pluginConfig); + String apiKey = pluginConfig.get(SplunkSourceOptions.API_KEY); + splunkSourceParameter.buildWithConfig(pluginConfig, apiKey); + } + + @Override + public String getPluginName() { + return "splunk"; + } + + @Override + public AbstractSingleSplitReader<SeaTunnelRow> createReader( + SingleSplitReaderContext readerContext) throws Exception { + return new HttpSourceReader( Review Comment: Confirmed against the current head — I had this flagged as unresolved from my first review but hadn't reproduced OOM myself, so thank you for the concrete repro. Traced the full buffering chain: `HttpClientProvider.execute()` materializes the entire response body as a `String` via `EntityUtils.toString(...)` (`HttpClientProvider.java:632`), then `SplunkSourceReader.executeRequest()` (`SplunkSourceReader.java:52-56`) hands that whole string to `filterAndUnwrapNdjson()`, which builds a second full copy via `StringBuilder` (`SplunkSourceReader.java:63-87`) before the parent's `pollAndCollectData()` (`HttpSourceReader.java:165-178`) does a third pass over it with `BufferedReader`/line splitting. There's no streaming boundary anywhere in that chain, so peak memory is a multiple of the raw response size — consistent with your 30MB-response/96MB-heap repro. This needs to become a single streaming pass (read the response body line-by-line and unwrap/emit per line, or use Splunk's paginated job/results flow) instead of three sequential full-body copies. Treating this as a High-severity blocker, not resolved by the current head — it's the same class of unbounded-materialization risk we just fixed for engine log responses in #12299, for what it's worth. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
