-
Notifications
You must be signed in to change notification settings - Fork 735
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Merge pull request #2945 from dantengsky/fix-2914
ISSUE-2814: introducing distributed insertion
- Loading branch information
Showing
38 changed files
with
982 additions
and
494 deletions.
There are no files selected for viewing
Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.
Oops, something went wrong.
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,52 @@ | ||
// Copyright 2021 Datafuse Labs. | ||
// | ||
// Licensed under the Apache License, Version 2.0 (the "License"); | ||
// you may not use this file except in compliance with the License. | ||
// You may obtain a copy of the License at | ||
// | ||
// http://www.apache.org/licenses/LICENSE-2.0 | ||
// | ||
// Unless required by applicable law or agreed to in writing, software | ||
// distributed under the License is distributed on an "AS IS" BASIS, | ||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. | ||
// See the License for the specific language governing permissions and | ||
// limitations under the License. | ||
// | ||
|
||
use common_base::tokio; | ||
use common_dal::DataAccessor; | ||
use common_dal::Local; | ||
use tempfile::TempDir; | ||
|
||
async fn local_read(loops: u32) -> common_exception::Result<()> { | ||
let tmp_root_dir = TempDir::new().unwrap(); | ||
let root_path = tmp_root_dir.path().to_str().unwrap(); | ||
let local_da = Local::new(root_path); | ||
|
||
let mut files = vec![]; | ||
for i in 0..loops { | ||
let file = format!("test_{}", i); | ||
let random_bytes: Vec<u8> = (0..122).map(|_| rand::random::<u8>()).collect(); | ||
local_da.put(file.as_str(), random_bytes).await?; | ||
files.push(file) | ||
} | ||
|
||
for x in files { | ||
local_da.read(x.as_str()).await?; | ||
} | ||
Ok(()) | ||
} | ||
|
||
// enable this if need to re-produce issue #2997 | ||
#[tokio::test] | ||
#[ignore] | ||
async fn test_da_local_hangs() -> common_exception::Result<()> { | ||
let read_fut = local_read(100); | ||
futures::executor::block_on(read_fut) | ||
} | ||
|
||
#[tokio::test] | ||
async fn test_da_local_normal() -> common_exception::Result<()> { | ||
let read_fut = local_read(1000); | ||
read_fut.await | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
|
@@ -14,3 +14,4 @@ | |
|
||
mod aws_s3; | ||
mod azure_blob; | ||
mod local; |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,45 @@ | ||
// Copyright 2020 Datafuse Labs. | ||
// | ||
// Licensed under the Apache License, Version 2.0 (the "License"); | ||
// you may not use this file except in compliance with the License. | ||
// You may obtain a copy of the License at | ||
// | ||
// http://www.apache.org/licenses/LICENSE-2.0 | ||
// | ||
// Unless required by applicable law or agreed to in writing, software | ||
// distributed under the License is distributed on an "AS IS" BASIS, | ||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. | ||
// See the License for the specific language governing permissions and | ||
// limitations under the License. | ||
|
||
use std::sync::Arc; | ||
|
||
use common_datavalues::DataField; | ||
use common_datavalues::DataSchemaRef; | ||
use common_datavalues::DataSchemaRefExt; | ||
use common_datavalues::DataType; | ||
use common_meta_types::TableInfo; | ||
use lazy_static::lazy_static; | ||
|
||
use crate::PlanNode; | ||
|
||
lazy_static! { | ||
pub static ref SINK_SCHEMA: DataSchemaRef = DataSchemaRefExt::create(vec![ | ||
DataField::new("seg_loc", DataType::String, false), | ||
DataField::new("seg_info", DataType::String, false), | ||
]); | ||
} | ||
|
||
#[derive(serde::Serialize, serde::Deserialize, Clone, Debug, PartialEq)] | ||
pub struct SinkPlan { | ||
pub table_info: TableInfo, | ||
pub input: Arc<PlanNode>, | ||
pub cast_needed: bool, | ||
} | ||
|
||
impl SinkPlan { | ||
/// Return sink schema | ||
pub fn schema(&self) -> DataSchemaRef { | ||
SINK_SCHEMA.clone() | ||
} | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.