Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 3 additions & 1 deletion fluss-gateway/src/backend/fake.rs
Original file line number Diff line number Diff line change
Expand Up @@ -216,7 +216,9 @@ fn write_table_info(
schema = schema.column(name, data_type);
}
if let Some(keys) = primary_key {
schema = schema.primary_key(keys.iter().copied());
schema = schema
.primary_key(keys.iter().copied())
.expect("valid fixture primary key");
}
let descriptor = TableDescriptor::builder()
.schema(schema.build().expect("valid fixture schema"))
Expand Down
5 changes: 4 additions & 1 deletion fluss-gateway/src/protocol/rest/ddl.rs
Original file line number Diff line number Diff line change
Expand Up @@ -520,7 +520,9 @@ fn table_descriptor(database: &str, body: CreateTableBody) -> GatewayResult<Tabl
}
}
if let Some(primary_key) = body.primary_key {
schema = schema.primary_key(primary_key);
schema = schema.primary_key(primary_key).map_err(|error| {
GatewayError::invalid_argument(format!("invalid table schema: {error}"))
})?;
}
let schema = schema.build().map_err(|error| {
GatewayError::invalid_argument(format!("invalid table schema: {error}"))
Expand Down Expand Up @@ -782,6 +784,7 @@ mod tests {
)
.with_comment("the order total")
.primary_key(["id", "dt"])
.unwrap()
.build()
.expect("a fixture schema");
let (status, location, body) = post(
Expand Down
1 change: 1 addition & 0 deletions fluss-gateway/src/protocol/rest/metadata.rs
Original file line number Diff line number Diff line change
Expand Up @@ -708,6 +708,7 @@ mod tests {
.with_comment("the order total")
.column("dt", DataType::String(StringType::with_nullable(false)))
.primary_key(["id", "dt"])
.unwrap()
.build()
.expect("the described schema is valid");
let descriptor = TableDescriptor::builder()
Expand Down
2 changes: 2 additions & 0 deletions fluss-gateway/src/protocol/rest/records.rs
Original file line number Diff line number Diff line change
Expand Up @@ -755,6 +755,7 @@ mod tests {
DataType::String(fluss::metadata::StringType::with_nullable(false)),
)
.primary_key(["id"])
.unwrap()
.build()
.unwrap();
let descriptor = TableDescriptor::builder()
Expand Down Expand Up @@ -787,6 +788,7 @@ mod tests {
DataType::BigInt(fluss::metadata::BigIntType::with_nullable(true)),
)
.primary_key(["uid"])
.unwrap()
.enable_auto_increment("uid_int")
.unwrap()
.build()
Expand Down
2 changes: 2 additions & 0 deletions fluss-gateway/tests/e2e_cluster.rs
Original file line number Diff line number Diff line change
Expand Up @@ -515,6 +515,7 @@ async fn assert_schema_recreation(api: &Api, connection: &FlussConnection) {
.column(columns[0], DataTypes::string())
.column(columns[1], DataTypes::string())
.primary_key(["id"])
.unwrap()
.build()
.unwrap(),
)
Expand Down Expand Up @@ -726,6 +727,7 @@ async fn metadata_apis_support_plaintext_and_sasl_fluss_clusters() {
.column("name", DataTypes::string())
.column("note", DataTypes::string())
.primary_key(["id"])
.unwrap()
.build()
.expect("build the KV schema"),
)
Expand Down
2 changes: 1 addition & 1 deletion fluss-rust/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -198,7 +198,7 @@ async fn main() -> Result<()> {
.column("id", DataTypes::int())
.column("name", DataTypes::string())
.column("score", DataTypes::bigint())
.primary_key(vec!["id"])
.primary_key(vec!["id"])?
.build()?;
let descriptor = TableDescriptor::builder().schema(schema).build()?;
admin.create_table(&table_path, &descriptor, true).await?;
Expand Down
2 changes: 1 addition & 1 deletion fluss-rust/bindings/cpp/src/types.rs
Original file line number Diff line number Diff line change
Expand Up @@ -244,7 +244,7 @@ pub fn ffi_descriptor_to_core(
}

if !descriptor.schema.primary_keys.is_empty() {
schema_builder = schema_builder.primary_key(descriptor.schema.primary_keys.clone());
schema_builder = schema_builder.primary_key(descriptor.schema.primary_keys.clone())?;
}

for auto_increment_column in &descriptor.schema.auto_increment_columns {
Expand Down
4 changes: 3 additions & 1 deletion fluss-rust/bindings/elixir/native/fluss_nif/src/schema.rs
Original file line number Diff line number Diff line change
Expand Up @@ -178,7 +178,9 @@ fn table_descriptor_new(
schema_builder = schema_builder.column(name, to_fluss_type(dt));
}
if !schema.primary_key.is_empty() {
schema_builder = schema_builder.primary_key(schema.primary_key);
schema_builder = schema_builder
.primary_key(schema.primary_key)
.map_err(to_nif_err)?;
}

let built_schema = schema_builder.build().map_err(to_nif_err)?;
Expand Down
4 changes: 3 additions & 1 deletion fluss-rust/bindings/python/src/metadata.rs
Original file line number Diff line number Diff line change
Expand Up @@ -176,7 +176,9 @@ impl Schema {

if let Some(pk_columns) = primary_keys {
if !pk_columns.is_empty() {
builder = builder.primary_key(pk_columns);
builder = builder
.primary_key(pk_columns)
.map_err(|e| FlussError::new_err(format!("Failed to build schema: {e}")))?;
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,7 @@ pub async fn main() -> Result<()> {
.column("uid", DataTypes::string())
.column("region", DataTypes::string())
.column("uid_int", DataTypes::bigint())
.primary_key(vec!["uid"])
.primary_key(vec!["uid"])?
.enable_auto_increment("uid_int")?
.build()?,
)
Expand Down
2 changes: 1 addition & 1 deletion fluss-rust/crates/examples/src/example_kv_changelog.rs
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,7 @@ pub async fn main() -> Result<()> {
Schema::builder()
.column("id", DataTypes::int())
.column("name", DataTypes::string())
.primary_key(vec!["id"])
.primary_key(vec!["id"])?
.build()?,
)
.distributed_by(Some(1), vec!["id".to_string()])
Expand Down
2 changes: 1 addition & 1 deletion fluss-rust/crates/examples/src/example_kv_table.rs
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,7 @@ pub async fn main() -> Result<()> {
.column("id", DataTypes::int())
.column("name", DataTypes::string())
.column("age", DataTypes::bigint())
.primary_key(vec!["id"])
.primary_key(vec!["id"])?
.build()?,
)
.build()?;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,7 @@ pub async fn main() -> Result<()> {
.column("region", DataTypes::string())
.column("zone", DataTypes::bigint())
.column("score", DataTypes::bigint())
.primary_key(vec!["id", "region", "zone"])
.primary_key(vec!["id", "region", "zone"])?
.build()?,
)
.partitioned_by(vec!["region", "zone"])
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,7 @@ pub async fn main() -> Result<()> {
.column("session_id", DataTypes::string())
.column("event_seq", DataTypes::bigint())
.column("event_data", DataTypes::string())
.primary_key(vec!["region", "user_id", "session_id", "event_seq"])
.primary_key(vec!["region", "user_id", "session_id", "event_seq"])?
.build()?,
)
.partitioned_by(vec!["region"])
Expand Down
2 changes: 1 addition & 1 deletion fluss-rust/crates/examples/src/example_prefix_lookup.rs
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,7 @@ pub async fn main() -> Result<()> {
.column("session_id", DataTypes::string())
.column("event_seq", DataTypes::bigint())
.column("event_data", DataTypes::string())
.primary_key(vec!["user_id", "session_id", "event_seq"])
.primary_key(vec!["user_id", "session_id", "event_seq"])?
.build()?,
)
.distributed_by(
Expand Down
4 changes: 2 additions & 2 deletions fluss-rust/crates/fluss/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -27,11 +27,11 @@ async fn main() -> Result<()> {

// ---- Primary key (KV) table: upsert and lookup ----
let kv_path = TablePath::new("fluss", "users");
let mut kv_schema = Schema::builder()
let kv_schema = Schema::builder()
.column("id", DataTypes::int())
.column("name", DataTypes::string())
.column("age", DataTypes::bigint())
.primary_key(vec!["id"]);
.primary_key(vec!["id"])?;
let kv_descriptor = TableDescriptor::builder()
.schema(kv_schema.build()?)
.build()?;
Expand Down
2 changes: 1 addition & 1 deletion fluss-rust/crates/fluss/src/client/table/scanner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3140,7 +3140,7 @@ mod tests {
.column("name", DataTypes::string());

if has_primary_key {
schema_builder = schema_builder.primary_key(vec!["id"]);
schema_builder = schema_builder.primary_key(vec!["id"]).unwrap();
}

let schema = schema_builder.build().unwrap();
Expand Down
4 changes: 2 additions & 2 deletions fluss-rust/crates/fluss/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -44,11 +44,11 @@
//!
//! // ---- Primary key (KV) table: upsert and lookup ----
//! let kv_path = TablePath::new("fluss", "users");
//! let mut kv_schema = Schema::builder()
//! let kv_schema = Schema::builder()
//! .column("id", DataTypes::int())
//! .column("name", DataTypes::string())
//! .column("age", DataTypes::bigint())
//! .primary_key(vec!["id"]);
//! .primary_key(vec!["id"])?;
//! let kv_descriptor = TableDescriptor::builder()
//! .schema(kv_schema.build()?)
//! .build()?;
Expand Down
4 changes: 3 additions & 1 deletion fluss-rust/crates/fluss/src/metadata/json_serde.rs
Original file line number Diff line number Diff line change
Expand Up @@ -527,7 +527,7 @@ impl JsonSerde for Schema {
);
}

schema_builder = schema_builder.primary_key(primary_keys);
schema_builder = schema_builder.primary_key(primary_keys)?;
}

if let Some(auto_increment_node) = node.get(Self::AUTO_INCREMENT_COLUMN_NAME) {
Expand Down Expand Up @@ -763,6 +763,7 @@ mod tests {
.column("id", DataTypes::int())
.column("seq", DataTypes::bigint())
.primary_key(["id"])
.unwrap()
.enable_auto_increment("seq")
.unwrap()
.build()
Expand All @@ -789,6 +790,7 @@ mod tests {
let schema = Schema::builder()
.column("id", DataTypes::int())
.primary_key(["id"])
.unwrap()
.build()
.unwrap();

Expand Down
1 change: 1 addition & 0 deletions fluss-rust/crates/fluss/src/metadata/schema_util.rs
Original file line number Diff line number Diff line change
Expand Up @@ -196,6 +196,7 @@ mod tests {
let origin = Schema::builder()
.column("a", DataTypes::int())
.primary_key(["a"])
.unwrap()
.build()
.unwrap();
let expected = schema_with_ids(&[(0, "a", DataTypes::int())]);
Expand Down
71 changes: 68 additions & 3 deletions fluss-rust/crates/fluss/src/metadata/table.rs
Original file line number Diff line number Diff line change
Expand Up @@ -291,7 +291,7 @@ impl SchemaBuilder {
self
}

pub fn primary_key<I, S>(self, column_names: I) -> Self
pub fn primary_key<I, S>(self, column_names: I) -> Result<Self>
where
I: IntoIterator<Item = S>,
S: Into<String>,
Expand All @@ -307,12 +307,18 @@ impl SchemaBuilder {
mut self,
constraint_name: N,
column_names: Vec<P>,
) -> Self {
) -> Result<Self> {
if self.primary_key.is_some() {
return Err(IllegalArgument {
message: "Multiple primary keys are not supported.".to_string(),
});
}

self.primary_key = Some(PrimaryKey::new(
constraint_name.into(),
column_names.into_iter().map(|s| s.into()).collect(),
));
self
Ok(self)
}

/// Declares a column to be auto-incremented. With an auto-increment column in the table,
Expand Down Expand Up @@ -531,6 +537,19 @@ impl SchemaBuilder {
return Ok(columns.to_vec());
};

if pk.column_names.is_empty() {
return Err(Error::invalid_table(
"Primary key constraint must be defined for at least a single column.",
));
}

let primary_key_names: Vec<_> = pk.column_names.iter().collect();
if let Some(duplicates) = Self::find_duplicates(&primary_key_names) {
return Err(Error::invalid_table(format!(
"Primary key constraint must not contain duplicate columns. Found: {duplicates:?}"
)));
}

let pk_set: HashSet<_> = pk.column_names.iter().collect();
let all_columns: HashSet<_> = columns.iter().map(|c| &c.name).collect();
if !pk_set.is_subset(&all_columns) {
Expand Down Expand Up @@ -1652,6 +1671,48 @@ mod tests {
use super::*;
use crate::metadata::DataTypes;

#[test]
fn invalid_primary_keys_are_rejected() {
for (primary_keys, expected_message) in [
(
Vec::<&str>::new(),
"Primary key constraint must be defined for at least a single column.",
),
(
vec!["id", "id"],
"Primary key constraint must not contain duplicate columns.",
),
] {
let err = Schema::builder()
.column("id", DataTypes::int())
.primary_key(primary_keys)
.unwrap()
.build()
.unwrap_err();

assert!(
err.to_string().contains(expected_message),
"unexpected error: {err}"
);
}
}

#[test]
fn multiple_primary_keys_are_rejected() {
let err = Schema::builder()
.column("id", DataTypes::int())
.primary_key(["id"])
.unwrap()
.primary_key_named("another_pk", vec!["id"])
.unwrap_err();

assert!(
err.to_string()
.contains("Multiple primary keys are not supported."),
"unexpected error: {err}"
);
}

#[test]
fn auto_increment_column_requires_a_primary_key_table() {
let err = Schema::builder()
Expand All @@ -1674,6 +1735,7 @@ mod tests {
.column("id", DataTypes::bigint())
.column("name", DataTypes::string())
.primary_key(["id"])
.unwrap()
.enable_auto_increment("id")
.unwrap()
.build()
Expand All @@ -1691,6 +1753,7 @@ mod tests {
.column("id", DataTypes::int())
.column("seq", DataTypes::string())
.primary_key(["id"])
.unwrap()
.enable_auto_increment("seq")
.unwrap()
.build()
Expand All @@ -1706,6 +1769,7 @@ mod tests {
.column("id", DataTypes::int())
.column("seq", accepted)
.primary_key(["id"])
.unwrap()
.enable_auto_increment("seq")
.unwrap()
.build()
Expand Down Expand Up @@ -1771,6 +1835,7 @@ mod tests {
.column("id", DataTypes::int())
.column("name", DataTypes::string())
.primary_key(vec!["id".to_string()])
.unwrap()
.build()
.unwrap();

Expand Down
3 changes: 3 additions & 0 deletions fluss-rust/crates/fluss/tests/integration/admin.rs
Original file line number Diff line number Diff line change
Expand Up @@ -99,6 +99,7 @@ mod admin_test {
.with_comment("User's age (optional)")
.column("email", DataTypes::string())
.primary_key(vec!["id".to_string()])
.unwrap()
.build()
.expect("Failed to build table schema");

Expand Down Expand Up @@ -217,6 +218,7 @@ mod admin_test {
.column("name", DataTypes::string())
.column("seq", DataTypes::bigint())
.primary_key(vec!["id".to_string()])
.unwrap()
.enable_auto_increment("seq")
.expect("Failed to enable auto increment")
.build()
Expand Down Expand Up @@ -292,6 +294,7 @@ mod admin_test {
.column("dt", DataTypes::string())
.column("region", DataTypes::string())
.primary_key(vec!["id", "dt", "region"])
.unwrap()
.build()
.expect("Failed to build table schema");

Expand Down
1 change: 1 addition & 0 deletions fluss-rust/crates/fluss/tests/integration/batch_scanner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -363,6 +363,7 @@ mod batch_scanner_test {
.column("id", DataTypes::int())
.column("name", DataTypes::string())
.primary_key(vec!["id"])
.unwrap()
.build()
.expect("schema"),
)
Expand Down
Loading
Loading