risingwavelabs/risingwave · error · SinkError::Config

snowflake.stage is required for S3 writer

Error message

snowflake.stage is required for S3 writer

What it means

The S3-based Snowflake writer builds a COPY pipe that ingests from an S3 stage; `snowflake_task_context.stage` must be set. When it is None, pipe creation fails with this Config error because the stage name is needed for `build_create_pipe_sql`.

Solutions

  1. Add `snowflake.stage = '<stage_name>'` to the sink WITH options when using S3 mode
  2. Verify the stage exists in Snowflake (`CREATE STAGE ...`) and the role can read it
  3. If no stage is desired, disable S3 mode and use the JDBC writer

Example fix

-- before
WITH (connector='snowflake', snowflake.s3=true, ...)
-- after
WITH (connector='snowflake', snowflake.s3=true, snowflake.stage='MY_STAGE', ...)
Defensive patterns

Strategy: validation

Validate before calling

if opts.get("snowflake.s3").map_or(false, |v| v == "true") {
    assert!(opts.contains_key("snowflake.stage"), "snowflake.stage is required for S3 writer");
}

Prevention

When it happens

Trigger: Using a Snowflake sink with S3 mode enabled where the task context was built without a `snowflake.stage` option, then writing data triggers pipe creation at snowflake.rs:949.

Common situations: S3 mode enabled but the `snowflake.stage` option omitted; copying JDBC-mode sink definitions into S3 mode; stage deleted/renamed in Snowflake.

Understand the failure class

Background: "is required", "must be set", "missing required field": configuration validation errors across open-source libraries — this error's family across 36 libraries.

Related errors


AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11). Data as JSON: /api/errors/5db0c5ced5ca9c68. Report an issue: GitHub.

Appendix: source

Thrown at src/connector/src/sink/snowflake_redshift/snowflake.rs:949

                .await?;
        }
        Ok(())
    }

    pub async fn execute_create_pipe(&self) -> Result<()> {
        if let Some(pipe_name) = &self.snowflake_task_context.pipe_name {
            let table_name =
                if let Some(table_name) = self.snowflake_task_context.cdc_table_name.as_ref() {
                    table_name
                } else {
                    &self.snowflake_task_context.target_table_name
                };
            let create_pipe_sql = build_create_pipe_sql(
                table_name,
                &self.snowflake_task_context.database,
                &self.snowflake_task_context.schema_name,
                self.snowflake_task_context.stage.as_ref().ok_or_else(|| {
                    SinkError::Config(anyhow!("snowflake.stage is required for S3 writer"))
                })?,
                pipe_name,
                &self.snowflake_task_context.target_table_name,
            );
            self.jdbc_client
                .execute_sql_sync(vec![create_pipe_sql])
                .await?;
        }
        Ok(())
    }

    pub async fn execute_drop_legacy_pipe(&self) -> Result<()> {
        // Older upsert sinks created this pipe and refreshed it from a local timer. Remove it
        // before starting the task so no asynchronous Snowpipe load can race with MERGE/DELETE.
        if self.snowflake_task_context.task_name.is_some()
            && self.snowflake_task_context.stage.is_some()
        {
            let pipe_name = format!("{}_pipe", self.snowflake_task_context.target_table_name);

View on GitHub (pinned to 6469eb736d)