Skip to main content
Need help generating SQL? Use Claude Code or Cursor with the RisingWave MCP server to generate and run SQL interactively. RisingWave offers several methods for data ingestion, each tailored to different use cases. This guide covers the main patterns to help you choose the best approach for your needs. For a detailed comparison of core objects like Source, Table, Materialized View, and Sink, see our guide on Source, Table, MV, and Sink.

Supported sources

Below is a complete list of source connectors in RisingWave. Click a connector name to see the SQL syntax, options, and a sample statement for connecting RisingWave to the connector. For information on supported data formats and encodings, and whether you need to use CREATE SOURCE or CREATE TABLE with each format, see Data formats and encoding options.

Continuous streaming ingestion

What it is: Real-time, continuous data ingestion from streaming sources that automatically updates as new data arrives. When to use: For real-time analytics, event-driven applications, live dashboards, and when you need immediate data freshness.

Option: HTTP ingestion (Webhook / Events API)

If you want to ingest events over HTTP without introducing Kafka or another message broker, you can use one of these options:
  • Webhook connector: RisingWave serves as the webhook destination and ingests requests into connector = 'webhook' tables (supports provider-style request validation/signatures). See Ingest data from webhook.
  • Events API: Run a standalone service that ingests JSON/NDJSON over HTTP and can execute SQL over HTTP. See Events API.

Example: Kafka streaming ingestion

Example: Database CDC (Change Data Capture)

Example: Message queues (MQTT, NATS, Pulsar)

One-time batch ingestion

What it is: Loading data once from external sources like databases, data lakes, or files. When to use: For initial data loads, historical data import, or when you need to load static datasets.

Example: Batch load from a database

For PostgreSQL and MySQL, use the postgres_query or mysql_query table-valued function to query the upstream database on demand. Unlike a CDC connector, each invocation executes the upstream query once and does not continue reading database changes.
CDC connectors are intended for continuous replication. Do not use a Debezium snapshot mode as a substitute for a finite batch load.

Example: Load from cloud storage (S3, GCS, Azure)

When a batch query such as CREATE TABLE AS or INSERT ... SELECT reads an object-storage source, RisingWave lists matching objects once for that query. In streaming queries, the same source continues discovering new files instead.

Example: Load from a data lake (Iceberg)

The batch query loads the current Iceberg snapshot once. To continuously ingest new snapshots from an append-only Iceberg table, use the source in a streaming query instead.

Periodic ingestion

What it is: Reloading or incrementally importing external data on a schedule. When to use: For scheduled data updates, daily or hourly batch processing, or when you need precise control over ingestion timing. S3, GCS, Iceberg, and Snowflake tables can use refresh_mode = 'FULL_RELOAD'. Set refresh_interval_sec for native periodic refresh, or omit it and run REFRESH TABLE when needed. Full-reload refresh is currently in technical preview. Use an external orchestrator such as Cron or Airflow for connectors without native refresh support or for custom incremental-loading logic.

Example: Externally orchestrated incremental load

A common pattern is to use a control table to track the last load time and only ingest new data.

Other ingestion methods

Direct data insertion

You can always insert data directly into a standard table using the INSERT statement.

Test data generation

For development and testing, you can use the built-in datagen connector to generate mock data streams.

Ingestion method support matrix

Legend:
  • ✅ Natively Supported: Built-in support for this ingestion method.
  • ❌ Not Supported: This ingestion method is not available for this source.
  • ⚠️ External Tools Required: Requires external orchestration tools (e.g., Cron, Airflow).

Best practices

Choose the right method

  • Streaming: Use for real-time requirements and continuous data flows.
  • Batch: Use for historical data, large one-time loads, or static datasets.
  • Periodic: Use native full-reload scheduling where supported, or external orchestration for custom incremental loads and other connectors.

Performance considerations

  • Streaming ingestion offers the best real-time performance.
  • Batch loading is efficient for large datasets.
  • Use materialized views to pre-compute and store results for fast querying.

Data consistency

  • CDC provides high-fidelity replication of database changes.
  • For message queues, understand the delivery guarantees (e.g., at-least-once) of your system.
  • Use transactions for atomic operations when inserting data manually.
  • Monitor data quality and set up alerts.

Monitoring and operations

  • Monitor streaming lag for real-time sources to ensure data freshness.
  • Track batch job success and failure rates.
  • Set up alerts for data quality issues.
  • Use RisingWave’s system tables and dashboards for monitoring.

See also