Configure integrations against real protocols.
Cloud services, analytical databases, discovery and telemetry boundaries aligned with FlowPort 0.1.5.
Catalog and configuration
41 native Go Source types and 34 Destination types, plus 2 Source and 8 Destination service presets. Presets reuse Kafka, S3 and Remote Write; they do not add transport protocols. Installation and runtime require neither Node.js nor a Vector process.
Select a service card and configure real endpoints, existing resources, authentication and protocol options. Use advanced JSON for complex rules. Drafts do not alter running versions; republish changes. Certificate paths and ambient identities belong to the actual execution Worker.
Service presets
CLS and Azure Event Hubs inputs reuse Kafka consumption; SLS, CLS and Event Hubs outputs reuse Kafka publishing. Use provider endpoints, Topics, consumer groups and SASL_SSL/PLAIN credentials. SLS publishing uses Project as username, AccessKey ID#Secret as password and Logstore.json as Topic. Event Hubs uses $ConnectionString as username and a namespace connection string as password. Kafka output presets disable idempotence and use gzip where required; Event Hubs consumption defaults to a 600ms commit interval.
R2 and MinIO outputs reuse S3. Configure the actual S3 Endpoint, Bucket, Access Key and Secret. R2 uses region=auto; MinIO defaults to us-east-1. Both presets use Path Style. A product name alone does not establish protocol or permission compatibility.
VictoriaMetrics single uses /api/v1/write; cluster vminsert uses /insert/<accountID[:projectID]>/prometheus/api/v1/write. Mimir uses /api/v1/push and tenant_id sets X-Scope-OrgID. Supply complete write URLs; tenant identity does not replace gateway authentication.
Native SLS Source
aliyun_sls requires endpoint, project, logstore, group_id, access_key and secret_key. This endpoint is the SLS API, not Kafka. Create the Project, Logstore and consumer group first. Use latest/earliest; polling defaults to one second.
Consumer heartbeats, shard checkpoints and takeover are supported. Newly acquired shards share a checkpoint snapshot. Commit only after local durable acceptance; capacity rejection or lease expiry does not acknowledge. Restart or takeover can duplicate records.
CloudWatch Logs
The aws_cloudwatch_logs Source needs region and log_group_name, with optional log_stream_name, filter_pattern and RFC3339 start_time. Default latest fixes the initial boundary; earliest supports backfill. New-event pages and watermarks persist atomically, and restart resumes from accepted progress.
lookback_seconds defaults to 60 (1–86400). Deduplicate event IDs within that window, retaining at most 100000 IDs. Distinct IDs at the same millisecond are accepted. Events older than the window require separate backfill. Duplicate pages are not enqueued; scan completion saves progress. An empty page does not end pagination.
The Destination uses PutLogEvents with region and existing log_group_name/log_stream_name. It accepts logs only and creates no resources. Source reads and Destination writes need different IAM permissions; diagnostics do not prove write permission.
Kinesis Data Streams
aws_kinesis Source and Destination need region and stream_name. This differs from Firehose HTTP ingestion. Sources discover shards, drain parents before children, restore iterators, verify/expand KPL records and persist sequence numbers. latest stores its initial time boundary. One Worker exclusively runs a flow.
Advance the physical KPL sequence only after all inner records are durable. Interruption can duplicate accepted inner records. Preserve Stream, Shard, Sequence, Subsequence and Partition Key tags.
Output uses PutRecords with at most 500 records/5 MiB per batch. This implementation limits each record plus key to 1 MiB, and partition keys to 1–256 bytes. Configure partition_key or partition_key_tag. Retry only transiently failed records within a call; queue retries may repeat successes. Global ordering and exactly-once delivery are not guaranteed.
RocketMQ
Requires an existing Topic, consumer group and RocketMQ 5.x gRPC Proxy at host:port. The 4.x Remoting protocol is not supported. Source uses SimpleConsumer with group_id, topic, filter and invisible_seconds (default 120, range 30–43200).
Renew message invisibility during processing and ACK after persistence. Capacity, lease and processing failures do not acknowledge. MESSAGE_NOT_FOUND is idle polling: retain the connection and pause 200ms. Destination publishes synchronously and checks receipts, with at most 4 MiB per message. Real RocketMQ 5.3.3 tests cover publishing after idle polling, consumption, durable acceptance and send receipts.
Doris and StarRocks
Choose doris/starrocks and configure url, database and table, with Basic authentication as required. Use Stream Load JSON with optional columns, jsonpaths, where, timezone and gzip.
Stable labels derive from queue batch identity and data. Validate business status, loaded rows and filtered rows. An existing label succeeds only when its original job is FINISHED. Partial writes are permanent failures; written rows are not automatically resent.
307/308 redirects require the same scheme, no embedded credentials and an explicitly allowed redirect_hosts host:port. Credentials are not forwarded to unapproved hosts. Diagnostics do not import business data; validate the real table and write path separately.
GCS and Azure object notifications
GCS Source needs project, subscription and bucket. Read OBJECT_FINALIZE notifications through an existing Pub/Sub subscription, fetch only the configured Bucket and specific generation, then acknowledge after persistence. Use Worker default GCP identity or credentials_json, with separate subscription and object-read permissions.
Azure Source needs url, container, queue and servicebus_namespace or connection_string. Route Event Grid BlobCreated events to an existing Service Bus queue. Fetch only the configured account/container with matching ETag, renew locks and complete after persistence. Blob supports Entra ID, Shared Key or SAS in separate fields. Blob SAS does not replace Service Bus identity.
Support gzip, chunked NDJSON/text and bounded JSON. Decompressed JSON is limited to 32 MiB; an NDJSON line to 1 MiB. No local-directory or whole-Bucket scan is performed. GCS/Azure Destinations write objects; reading notifications and writing objects are different roles.
Prometheus discovery
Combine static url/urls, HTTP SD, DNS SRV/A/AAAA and remote Kubernetes EndpointSlices/Consul passing services. Configure kubernetes_sd_url or consul_sd_url with separate provider credentials. HTTP SD shares authentication/TLS with scraping; trust the discovery service.
Advanced JSON supports target_relabel_configs and metric_relabel_configs (up to 64 each), plus full-name metric_include/metric_exclude regular expressions, one per line. Parse rules once per running version; republish changes.
Deduplicate by the relabeled URL and ordinary labels. The same URL with distinct labels remains separate; identical targets are scraped once. max_targets defaults to 256 (1–1000), concurrency to 4 (1–32). Failed discovery retains previous process-local targets; successful empty results remove them. Caches do not survive restart.
Telemetry fidelity and output boundaries
OTLP HTTP/gRPC retains Resource, Scope, schema, Sum temporality/monotonic, Histogram, ExponentialHistogram, Summary, exemplars and Span Kind/status/events/Links. Merge Pipeline numeric/tag edits into output; dropped Points remove corresponding protocol records.
Remote Write v1 input retains native histograms and exemplars. Ordinary Points and other protocols do not receive automatic equivalent native-histogram conversion. Metadata and Remote Write v2 are not emitted. Ordinary samples still use float64 and milliseconds, without arbitrary integer or nanosecond precision.
Vector gRPC uses v2 PushEvents. Ordinary Point metrics map to Gauges; original Counter/Histogram types and Vector envelopes are not restored. Internal __flowport_ fields are hidden in JSON, text and Line Protocol output. Partial OTLP rejection is not full success or automatically resent. Validate actual downstream receipts.
Diagnostics and validation scope
Test connections and sampling in the actual Worker. Diagnostics do not prove write permission, table compatibility or business delivery. Messaging/push and continuous cloud Sources sample running input without adding consumers or advancing cursors.
This release passed Go tests, race checks, protocol fixtures, real RocketMQ receipt tests, five-language checks and Node.js-free offline installation/runtime checks on both Linux architectures. Validate real cloud IAM, quotas, networks and sustained load in the deployment environment. Community has one local Worker; multiple Workers need an offline license.