DuckDB长时间运行时无法同步S3 Iceberg表?咨询原因
DuckDB挂载S3 Iceberg表后无法获取新增数据的问题解决
问题场景
已通过DuckDB成功挂载S3上的Iceberg表,挂载代码如下:
query = fmt.Sprintf(` ATTACH '%s' AS s3_tables ( TYPE ICEBERG, ENDPOINT_TYPE s3_tables );`, config.GetAwsS3TableArn(), ) _, err = r.duckdb.ExecContext(ctx, query) if err != nil { return fmt.Errorf("failed to attach s3_tables: %v", err) }
但在长时间运行的进程中执行两次查询(间隔10分钟,期间Firehose每5分钟向Iceberg表插入新数据),两次返回结果完全一致,无法获取新增数据:
func (m *MockApp) TryReadFromS3TablesDirectly(ctx context.Context) error { entries, err := m.repos.DuckDBSensorEntry.PreviewLatestData(ctx) if err != nil { return fmt.Errorf("failed to get first preview latest data: %v", err) } for _, e := range entries { m.logger.Info().Any("entry", e).Msg("preview first data") } time.Sleep(10 * time.Minute) entries, err = m.repos.DuckDBSensorEntry.PreviewLatestData(ctx) if err != nil { return fmt.Errorf("failed to get second preview latest data: %v", err) } for _, e := range entries { m.logger.Info().Any("entry", e).Msg("preview second data") } return nil }
原因说明
这不是Iceberg本身的限制,Iceberg支持实时更新元数据和数据文件。问题出在DuckDB的默认行为:挂载Iceberg表后,DuckDB会缓存表的元数据(包括快照、数据文件列表等),长时间运行的进程不会主动触发元数据刷新,因此后续查询依然使用旧的缓存数据,看不到新增内容。
解决办法
1. 每次查询前手动刷新元数据
在执行查询前,先执行REFRESH TABLE命令强制刷新Iceberg表的元数据,确保获取最新的快照和数据文件。修改PreviewLatestData方法,添加元数据刷新逻辑:
func (d *DuckDBSensorEntry) PreviewLatestData(ctx context.Context) ([]SensorEntry, error) { // 刷新指定Iceberg表的元数据 _, err := d.db.ExecContext(ctx, "REFRESH TABLE s3_tables.your_sensor_table;") if err != nil { return nil, fmt.Errorf("failed to refresh iceberg table: %v", err) } // 执行原有查询逻辑 rows, err := d.db.QueryContext(ctx, "SELECT * FROM s3_tables.your_sensor_table LIMIT 10;") if err != nil { return nil, fmt.Errorf("failed to query sensor data: %v", err) } defer rows.Close() // 解析结果... }
2. 开启自动元数据刷新
可以通过设置DuckDB配置参数,让每次查询自动刷新Iceberg表的元数据。注意该设置会增加查询的额外开销,需根据业务场景权衡:
// 在挂载表后或初始化时设置参数 _, err = r.duckdb.ExecContext(ctx, "SET iceberg_refresh_metadata_on_query = true;") if err != nil { return fmt.Errorf("failed to set iceberg refresh config: %v", err) }
3. 重新挂载表(不推荐)
如果上述两种方式不适用,也可以在每次查询前重新挂载表,但重新挂载的开销远大于元数据刷新,仅在特殊场景下使用。
内容的提问来源于stack exchange,提问作者Muhammad Najid
相关产品推荐
相关产品推荐

