HTTP Pull Source Connector
stream sourcescan table sourcelookup table source
The HTTP Pull source connector retrieves data periodically from external HTTP servers.
The connector supports fixed polling intervals, conditional fetching, dynamic URL templating, and OAuth 2.0 token management.
Configuration Overview
Configure the connector using environment variables, the REST API, or the configuration file.
The configuration file resides at $rekuiper/etc/sources/http_pull.yaml. Properties defined in the default section provide global default values. Custom sections override default values.
Example configuration file:
default:
url: http://localhost
method: post
interval: 10000
timeout: 5000
incremental: false
body: '{}'
bodyType: json
insecureSkipVerify: true
headers:
Accept: application/json
states:
responseType: code
application_conf:
incremental: true
url: http://localhost:9090/pullGlobal Configurations
Properties in the default section apply to all HTTP connections unless explicitly overridden.
HTTP Request Parameters
url: Target endpoint URL.method: HTTP method:post,get,put, ordelete.interval: Polling interval in milliseconds.timeout: Request timeout in milliseconds.body: Request payload string (for example,'{"data": "telemetry"}').bodyType: Payload content format:none,text,json,html,xml,javascript, orform.headers: Map of HTTP headers sent with the request.states: Key-value pairs rendered into URL template parameters.responseType: Response validation method:code: Validates success based on the HTTP response status code.body: Validates success based on acodefield inside a JSON response body.
TLS and Security Parameters
certificationPath: Path to client certificate file in PEM format.privateKeyPath: Path to client private key file in PEM format.rootCaPath: Path to Root CA certificate file in PEM format.certficationRaw: Base64-encoded client certificate string.privateKeyRaw: Base64-encoded client private key string.rootCARaw: Base64-encoded Root CA certificate string.insecureSkipVerify: Boolean. Set totrueto skip certificate validation.
OAuth 2.0 Authentication
Configure OAuth 2.0 token acquisition and refresh under the oAuth property:
oAuth:
access:
url: https://127.0.0.1/api/token
body: '{"username": "admin", "password": "password"}'
expire: "3600"
refresh:
url: https://127.0.0.1/api/refresh
headers:
Accept: application/json
identityId: '{{.data.identityId}}'
token: '{{.data.token}}'
body: ''access.url: Endpoint URL for initial token retrieval using HTTP POST.access.body: Request payload containing user credentials or authorization codes.access.expire: Token lifetime in seconds. Supports template expressions.refresh.url: Endpoint URL for token refresh using HTTP POST.refresh.headers: Headers required for token refresh. Supports template values extracted from initial responses.refresh.body: Payload required for token refresh.
Incremental Fetching and State Tracking
Incremental Processing
Set incremental: true to prevent duplicate processing. The engine compares each response against the previous response. If values are identical, the connector discards the result.
Dynamic State Tracking
Configure the states property to store runtime values between polling cycles. The engine renders state values into URL parameters using Data Template syntax.
When qos: 1 is configured, the engine flushes state to disk during checkpoints and restores state upon restart.
Configuration example:
default:
url: http://localhost/path?key1={{.key1}}&key2={{.key2}}
method: get
interval: 10000
timeout: 5000
body: '{}'
bodyType: json
headers:
Accept: application/json
states:
key1: value1
key2: value2
responseType: codeThe first request fetches:
GET http://localhost/path?key1=value1&key2=value2If the response returns:
{
"key1": "value3",
"key2": "value4"
}The next request uses the updated values:
GET http://localhost/path?key1=value3&key2=value4Dynamic Properties
Dynamic properties resolve values at runtime:
PullTime: Unix timestamp of the current poll execution in milliseconds (int64).LastPullTime: Unix timestamp of the previous poll execution in milliseconds (int64).- OAuth response fields: Access tokens or session IDs returned by authentication endpoints.
Use these variables to query time-windowed external APIs:
- In URL parameters:
http://localhost:9090/pull?start={{.LastPullTime}}&end={{.PullTime}} - In request payloads:
{"start": {{.LastPullTime}}, "end": {{.PullTime}}}
Custom Configurations
Define named configuration blocks in http_pull.yaml:
application_conf:
incremental: true
url: http://localhost:9090/pullReference the configuration with CONF_KEY="application_conf" in the stream definition:
CREATE STREAM demo () WITH (
DATASOURCE = "test/",
FORMAT = "JSON",
TYPE = "httppull",
CONF_KEY = "application_conf"
);Create a Stream Source
The HTTP Pull connector operates as a stream source or as a scan table.
Create Stream via REST API
POST http://{{host}}/streams
Content-Type: application/json
{
"sql": "CREATE STREAM http_stream () WITH (FORMAT = \"json\", TYPE = \"httppull\");"
}Create Stream via CLI
bin/kuiper create stream http_stream '() WITH (FORMAT = "json", TYPE = "httppull")'Create a Lookup Table Source
HTTP Pull supports lookup tables. The engine executes on-demand HTTP requests when a rule joins a stream with the table:
CREATE TABLE httppullTable () WITH (
DATASOURCE = "/url",
CONF_KEY = "default",
TYPE = "httppull",
KIND = "lookup"
);