Skip to content

Commit

Permalink
Merge branch 'main' of https://github.com/risingwavelabs/risingwave i…
Browse files Browse the repository at this point in the history
…nto wcy/change_opendal_write_concurrency.pr
  • Loading branch information
wcy-fdu committed Sep 19, 2024
2 parents aca8453 + c7b6fff commit b6a03cb
Show file tree
Hide file tree
Showing 25 changed files with 176 additions and 325 deletions.
2 changes: 1 addition & 1 deletion ci/docker-compose.yml
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,7 @@ services:
retries: 5

message_queue:
image: "docker.vectorized.io/vectorized/redpanda:latest"
image: "redpandadata/redpanda:latest"
command:
- redpanda
- start
Expand Down
2 changes: 1 addition & 1 deletion docker/docker-compose-distributed-etcd.yml
Original file line number Diff line number Diff line change
Expand Up @@ -329,7 +329,7 @@ services:
retries: 5
restart: always
message_queue:
image: "docker.vectorized.io/vectorized/redpanda:latest"
image: "redpandadata/redpanda:latest"
command:
- redpanda
- start
Expand Down
2 changes: 1 addition & 1 deletion docker/docker-compose-distributed.yml
Original file line number Diff line number Diff line change
Expand Up @@ -299,7 +299,7 @@ services:
retries: 5
restart: always
message_queue:
image: "docker.vectorized.io/vectorized/redpanda:latest"
image: "redpandadata/redpanda:latest"
command:
- redpanda
- start
Expand Down
2 changes: 1 addition & 1 deletion docker/docker-compose-etcd.yml
Original file line number Diff line number Diff line change
Expand Up @@ -228,7 +228,7 @@ services:
restart: always

message_queue:
image: "docker.vectorized.io/vectorized/redpanda:latest"
image: "redpandadata/redpanda:latest"
command:
- redpanda
- start
Expand Down
2 changes: 1 addition & 1 deletion docker/docker-compose-with-hdfs.yml
Original file line number Diff line number Diff line change
Expand Up @@ -275,7 +275,7 @@ services:
retries: 5
restart: always
message_queue:
image: "docker.vectorized.io/vectorized/redpanda:latest"
image: "redpandadata/redpanda:latest"
command:
- redpanda
- start
Expand Down
2 changes: 1 addition & 1 deletion docker/docker-compose-with-sqlite.yml
Original file line number Diff line number Diff line change
Expand Up @@ -177,7 +177,7 @@ services:
restart: always

message_queue:
image: "docker.vectorized.io/vectorized/redpanda:latest"
image: "redpandadata/redpanda:latest"
command:
- redpanda
- start
Expand Down
2 changes: 1 addition & 1 deletion docker/docker-compose.yml
Original file line number Diff line number Diff line change
Expand Up @@ -199,7 +199,7 @@ services:
restart: always

message_queue:
image: "docker.vectorized.io/vectorized/redpanda:latest"
image: "redpandadata/redpanda:latest"
command:
- redpanda
- start
Expand Down
7 changes: 7 additions & 0 deletions e2e_test/batch/catalog/rw_depend.slt.part
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,9 @@ create materialized view mv1 as select t1.a from t1 join s1 on t1.a = s1.a;
statement ok
create materialized view mv2 as select * from mv1;

statement ok
create view v as select * from mv1;

statement ok
create sink sink1 from mv2 with (connector='blackhole');

Expand All @@ -37,13 +40,17 @@ mv2 mv1
sink1 mv2
sink2 t1
t2 sink2
v mv1

statement ok
drop sink sink1;

statement ok
drop table t2 cascade;

statement ok
drop view v;

statement ok
drop materialized view mv2;

Expand Down
21 changes: 4 additions & 17 deletions proto/hummock.proto
Original file line number Diff line number Diff line change
Expand Up @@ -87,21 +87,6 @@ message GroupConstruct {
CompatibilityVersion version = 6;
}

message GroupMetaChange {
option deprecated = true;
repeated uint32 table_ids_add = 1 [deprecated = true];
repeated uint32 table_ids_remove = 2 [deprecated = true];
}

message GroupTableChange {
option deprecated = true;
repeated uint32 table_ids = 1 [deprecated = true];
uint64 target_group_id = 2;
uint64 origin_group_id = 3;
uint64 new_sst_start_id = 4;
CompatibilityVersion version = 5;
}

message GroupDestroy {}

message GroupMerge {
Expand All @@ -110,12 +95,14 @@ message GroupMerge {
}

message GroupDelta {
reserved 4;
reserved "group_meta_change";
reserved 5;
reserved "group_table_change";
oneof delta_type {
IntraLevelDelta intra_level = 1;
GroupConstruct group_construct = 2;
GroupDestroy group_destroy = 3;
GroupMetaChange group_meta_change = 4 [deprecated = true];
GroupTableChange group_table_change = 5 [deprecated = true];
GroupMerge group_merge = 6;
}
}
Expand Down
14 changes: 5 additions & 9 deletions proto/stream_plan.proto
Original file line number Diff line number Diff line change
Expand Up @@ -463,22 +463,18 @@ message AsOfJoinNode {
catalog.Table left_table = 4;
// Used for internal table states.
catalog.Table right_table = 5;
// Used for internal table states.
catalog.Table left_degree_table = 6;
// Used for internal table states.
catalog.Table right_degree_table = 7;
// The output indices of current node
repeated uint32 output_indices = 8;
repeated uint32 output_indices = 6;
// Left deduped input pk indices. The pk of the left_table and
// The pk of the left_table is [left_join_key | left_inequality_key | left_deduped_input_pk_indices]
// left_inequality_key is not used but for forward compatibility.
repeated uint32 left_deduped_input_pk_indices = 9;
repeated uint32 left_deduped_input_pk_indices = 7;
// Right deduped input pk indices.
// The pk of the right_table is [right_join_key | right_inequality_key | right_deduped_input_pk_indices]
// right_inequality_key is not used but for forward compatibility.
repeated uint32 right_deduped_input_pk_indices = 10;
repeated bool null_safe = 11;
optional plan_common.AsOfJoinDesc asof_desc = 12;
repeated uint32 right_deduped_input_pk_indices = 8;
repeated bool null_safe = 9;
optional plan_common.AsOfJoinDesc asof_desc = 10;
}

message TemporalJoinNode {
Expand Down
5 changes: 4 additions & 1 deletion risedev.yml
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,9 @@ profile:
# - use: kafka
# persist-data: true

# To enable Confluent schema registry, uncomment the following line
# - use: schema-registry

default-v6:
steps:
- use: meta-node
Expand Down Expand Up @@ -1141,7 +1144,7 @@ compose:
risingwave: "ghcr.io/risingwavelabs/risingwave:latest"
prometheus: "prom/prometheus:latest"
minio: "quay.io/minio/minio:latest"
redpanda: "docker.vectorized.io/vectorized/redpanda:latest"
redpanda: "redpandadata/redpanda:latest"
grafana: "grafana/grafana-oss:latest"
etcd: "quay.io/coreos/etcd:latest"
tempo: "grafana/tempo:latest"
Expand Down
22 changes: 11 additions & 11 deletions src/connector/src/sink/encoder/avro.rs
Original file line number Diff line number Diff line change
Expand Up @@ -652,7 +652,7 @@ mod tests {
}
]
}"#,
"encode q error: avro name ref unsupported yet",
"encode 'q' error: avro name ref unsupported yet",
);

test_err(
Expand All @@ -663,7 +663,7 @@ mod tests {
i64::MAX,
))),
r#"{"type": "fixed", "name": "Duration", "size": 12, "logicalType": "duration"}"#,
"encode error: -1 mons -1 days +2562047788:00:54.775807 overflows avro duration",
"encode '' error: -1 mons -1 days +2562047788:00:54.775807 overflows avro duration",
);

let avro_schema = AvroSchema::parse_str(
Expand Down Expand Up @@ -738,7 +738,7 @@ mod tests {
};
assert_eq!(
err.to_string(),
"Encode error: encode req error: field not present but required"
"Encode error: encode 'req' error: field not present but required"
);

let schema = Schema::new(vec![
Expand All @@ -751,7 +751,7 @@ mod tests {
};
assert_eq!(
err.to_string(),
"Encode error: encode extra error: field not in avro"
"Encode error: encode 'extra' error: field not in avro"
);

let avro_schema = AvroSchema::parse_str(r#"["null", "long"]"#).unwrap();
Expand All @@ -761,14 +761,14 @@ mod tests {
};
assert_eq!(
err.to_string(),
r#"Encode error: encode error: expect avro record but got ["null","long"]"#
r#"Encode error: encode '' error: expect avro record but got ["null","long"]"#
);

test_err(
&DataType::Struct(StructType::new(vec![("f0", DataType::Boolean)])),
(),
r#"{"type": "record", "name": "T", "fields": [{"name": "f0", "type": "int"}]}"#,
"encode f0 error: cannot encode boolean column as \"int\" field",
"encode 'f0' error: cannot encode boolean column as \"int\" field",
);
}

Expand All @@ -790,7 +790,7 @@ mod tests {
&DataType::List(DataType::Int32.into()),
Some(ScalarImpl::List(ListValue::from_iter([Some(4), None]))).to_datum_ref(),
avro_schema,
"encode error: found null but required",
"encode '' error: found null but required",
);

test_ok(
Expand Down Expand Up @@ -829,7 +829,7 @@ mod tests {
&DataType::List(DataType::Boolean.into()),
(),
r#"{"type": "array", "items": "int"}"#,
"encode error: cannot encode boolean column as \"int\" field",
"encode '' error: cannot encode boolean column as \"int\" field",
);
}

Expand Down Expand Up @@ -863,14 +863,14 @@ mod tests {
t,
datum.to_datum_ref(),
both,
r#"encode error: cannot encode timestamp with time zone column as [{"type":"long","logicalType":"timestamp-millis"},{"type":"long","logicalType":"timestamp-micros"}] field"#,
r#"encode '' error: cannot encode timestamp with time zone column as [{"type":"long","logicalType":"timestamp-millis"},{"type":"long","logicalType":"timestamp-micros"}] field"#,
);

test_err(
t,
datum.to_datum_ref(),
empty,
"encode error: cannot encode timestamp with time zone column as [] field",
"encode '' error: cannot encode timestamp with time zone column as [] field",
);

test_ok(
Expand All @@ -879,7 +879,7 @@ mod tests {
one,
Value::Union(0, Value::TimestampMillis(1).into()),
);
test_err(t, None, one, "encode error: found null but required");
test_err(t, None, one, "encode '' error: found null but required");

test_ok(
t,
Expand Down
2 changes: 1 addition & 1 deletion src/connector/src/sink/encoder/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -201,7 +201,7 @@ impl std::fmt::Display for FieldEncodeError {

write!(
f,
"encode {} error: {}",
"encode '{}' error: {}",
self.rev_path.iter().rev().join("."),
self.message
)
Expand Down
8 changes: 4 additions & 4 deletions src/connector/src/sink/encoder/proto.rs
Original file line number Diff line number Diff line change
Expand Up @@ -529,7 +529,7 @@ mod tests {
.unwrap_err();
assert_eq!(
err.to_string(),
"encode repeated_int_field error: cannot encode integer[] column as int32 field"
"encode 'repeated_int_field' error: cannot encode integer[] column as int32 field"
);

let schema = Schema::new(vec![Field::with_name(
Expand All @@ -554,7 +554,7 @@ mod tests {
.unwrap_err();
assert_eq!(
err.to_string(),
"encode repeated_int_field error: array containing null not allowed as repeated field"
"encode 'repeated_int_field' error: array containing null not allowed as repeated field"
);
}

Expand All @@ -573,7 +573,7 @@ mod tests {
.unwrap_err();
assert_eq!(
err.to_string(),
"encode not_exists error: field not in proto"
"encode 'not_exists' error: field not in proto"
);

let err = validate_fields(
Expand All @@ -583,7 +583,7 @@ mod tests {
.unwrap_err();
assert_eq!(
err.to_string(),
"encode map_field error: field not in proto"
"encode 'map_field' error: field not in proto"
);
}
}
10 changes: 7 additions & 3 deletions src/connector/src/sink/formatter/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@
// See the License for the specific language governing permissions and
// limitations under the License.

use anyhow::anyhow;
use anyhow::{anyhow, Context};
use risingwave_common::array::StreamChunk;

use crate::sink::{Result, SinkError};
Expand Down Expand Up @@ -279,8 +279,12 @@ impl<KE: EncoderBuild, VE: EncoderBuild> FormatterBuild for AppendOnlyFormatter<

impl<KE: EncoderBuild, VE: EncoderBuild> FormatterBuild for UpsertFormatter<KE, VE> {
async fn build(b: FormatterParams<'_>) -> Result<Self> {
let key_encoder = KE::build(b.builder.clone(), Some(b.pk_indices)).await?;
let val_encoder = VE::build(b.builder, None).await?;
let key_encoder = KE::build(b.builder.clone(), Some(b.pk_indices))
.await
.with_context(|| "Failed to build key encoder")?;
let val_encoder = VE::build(b.builder, None)
.await
.with_context(|| "Failed to build value encoder")?;
Ok(UpsertFormatter::new(key_encoder, val_encoder))
}
}
Expand Down
2 changes: 1 addition & 1 deletion src/connector/src/sink/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -644,7 +644,7 @@ pub enum SinkError {
#[backtrace]
anyhow::Error,
),
#[error("Internal error: {0}")]
#[error(transparent)]
Internal(
#[from]
#[backtrace]
Expand Down
22 changes: 22 additions & 0 deletions src/meta/src/controller/catalog.rs
Original file line number Diff line number Diff line change
Expand Up @@ -674,6 +674,28 @@ impl CatalogController {
})
.collect_vec();

let view_dependencies: Vec<(ObjectId, ObjectId)> = ObjectDependency::find()
.select_only()
.columns([
object_dependency::Column::Oid,
object_dependency::Column::UsedBy,
])
.join(
JoinType::InnerJoin,
object_dependency::Relation::Object1.def(),
)
.join(JoinType::InnerJoin, object::Relation::View.def())
.into_tuple()
.all(&inner.db)
.await?;

obj_dependencies.extend(view_dependencies.into_iter().map(|(view_id, table_id)| {
PbObjectDependencies {
object_id: table_id as _,
referenced_object_id: view_id as _,
}
}));

let sink_dependencies: Vec<(SinkId, TableId)> = Sink::find()
.select_only()
.columns([sink::Column::SinkId, sink::Column::TargetTable])
Expand Down
8 changes: 8 additions & 0 deletions src/meta/src/manager/catalog/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4217,6 +4217,14 @@ impl CatalogManager {
}
}
}
for view in core.views.values() {
for referenced in &view.dependent_relations {
dependencies.push(PbObjectDependencies {
object_id: view.id,
referenced_object_id: *referenced,
});
}
}
for sink in core.sinks.values() {
for referenced in &sink.dependent_relations {
dependencies.push(PbObjectDependencies {
Expand Down
Loading

0 comments on commit b6a03cb

Please sign in to comment.