From 096b81b635d8f61571f7e9d1763e1784b217416f Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Fri, 11 Apr 2025 17:42:05 -0400 Subject: [PATCH 01/14] Remove archive metadata update from clp-s and add this to compression task --- .../scripts/native/compress.py | 7 +++++- .../clp-py-utils/clp_py_utils/clp_config.py | 1 + components/core/src/clp_s/ArchiveWriter.cpp | 24 ++----------------- components/core/src/clp_s/ArchiveWriter.hpp | 10 +------- components/core/src/clp_s/JsonParser.cpp | 2 +- .../executor/compress/compression_task.py | 17 +++++++++++++ 6 files changed, 28 insertions(+), 33 deletions(-) diff --git a/components/clp-package-utils/clp_package_utils/scripts/native/compress.py b/components/clp-package-utils/clp_package_utils/scripts/native/compress.py index b71907eb26..2af96be38e 100755 --- a/components/clp-package-utils/clp_package_utils/scripts/native/compress.py +++ b/components/clp-package-utils/clp_package_utils/scripts/native/compress.py @@ -61,6 +61,7 @@ def handle_job_update(db, db_cursor, job_id, no_progress_reporting): ) job_last_uncompressed_size = 0 + last_current_time = datetime.datetime.now() while True: db_cursor.execute(polling_query) @@ -79,7 +80,10 @@ def handle_job_update(db, db_cursor, job_id, no_progress_reporting): # All tasks in the job is done if not no_progress_reporting: logger.info("Compression finished.") - print_compression_job_status(job_row, current_time) + if job_last_uncompressed_size < job_uncompressed_size: + print_compression_job_status(job_row, current_time) + else: + print_compression_job_status(job_row, last_current_time) break # Done if CompressionJobStatus.FAILED == job_status: # One or more tasks in the job has failed @@ -98,6 +102,7 @@ def handle_job_update(db, db_cursor, job_id, no_progress_reporting): error_msg = f"Unhandled CompressionJobStatus: {job_status}" raise NotImplementedError(error_msg) + last_current_time = current_time time.sleep(0.5) diff --git a/components/clp-py-utils/clp_py_utils/clp_config.py b/components/clp-py-utils/clp_py_utils/clp_config.py index 1fbf5cbe63..b00c9cf673 100644 --- a/components/clp-py-utils/clp_py_utils/clp_config.py +++ b/components/clp-py-utils/clp_py_utils/clp_config.py @@ -36,6 +36,7 @@ QUERY_TASKS_TABLE_NAME = "query_tasks" COMPRESSION_JOBS_TABLE_NAME = "compression_jobs" COMPRESSION_TASKS_TABLE_NAME = "compression_tasks" +ARCHIVES_TABLE_NAME = "clp_archives" OS_RELEASE_FILE_PATH = pathlib.Path("etc") / "os-release" diff --git a/components/core/src/clp_s/ArchiveWriter.cpp b/components/core/src/clp_s/ArchiveWriter.cpp index df53d7e96d..a5e004b1cd 100644 --- a/components/core/src/clp_s/ArchiveWriter.cpp +++ b/components/core/src/clp_s/ArchiveWriter.cpp @@ -100,10 +100,6 @@ void ArchiveWriter::close() { header_and_metadata_writer.close(); } - if (m_metadata_db) { - update_metadata_db(); - } - if (m_print_archive_stats) { print_archive_stats(); } @@ -430,27 +426,11 @@ std::pair ArchiveWriter::store_tables() { return {table_metadata_compressed_size, table_compressed_size}; } -void ArchiveWriter::update_metadata_db() { - m_metadata_db->open(); - clp::streaming_archive::ArchiveMetadata metadata( - cArchiveFormatDevelopmentVersionFlag, - "", - 0ULL - ); - metadata.increment_static_compressed_size(m_compressed_size); - metadata.increment_static_uncompressed_size(m_uncompressed_size); - metadata.expand_time_range( - m_timestamp_dict.get_begin_timestamp(), - m_timestamp_dict.get_end_timestamp() - ); - - m_metadata_db->add_archive(m_id, metadata); - m_metadata_db->close(); -} - void ArchiveWriter::print_archive_stats() { nlohmann::json json_msg; json_msg["id"] = m_id; + json_msg["begin_timestamp"] = m_timestamp_dict.get_begin_timestamp(); + json_msg["end_timestamp"] = m_timestamp_dict.get_end_timestamp(); json_msg["uncompressed_size"] = m_uncompressed_size; json_msg["size"] = m_compressed_size; std::cout << json_msg.dump(-1, ' ', true, nlohmann::json::error_handler_t::ignore) << std::endl; diff --git a/components/core/src/clp_s/ArchiveWriter.hpp b/components/core/src/clp_s/ArchiveWriter.hpp index 7fb0be54cc..aa05a138fa 100644 --- a/components/core/src/clp_s/ArchiveWriter.hpp +++ b/components/core/src/clp_s/ArchiveWriter.hpp @@ -7,7 +7,6 @@ #include #include -#include "../clp/GlobalMySQLMetadataDB.hpp" #include "archive_constants.hpp" #include "DictionaryWriter.hpp" #include "Schema.hpp" @@ -66,8 +65,7 @@ class ArchiveWriter { }; // Constructor - explicit ArchiveWriter(std::shared_ptr metadata_db) - : m_metadata_db(std::move(metadata_db)) {} + ArchiveWriter() = default; // Destructor ~ArchiveWriter() = default; @@ -211,11 +209,6 @@ class ArchiveWriter { */ void write_archive_header(FileWriter& archive_writer, size_t metadata_section_size); - /** - * Updates the metadata db with the archive's metadata (id, size, timestamp ranges, etc.) - */ - void update_metadata_db(); - /** * Prints the archive's statistics (id, uncompressed size, compressed size, etc.) */ @@ -238,7 +231,6 @@ class ArchiveWriter { std::shared_ptr m_log_dict; std::shared_ptr m_array_dict; // log type dictionary for arrays TimestampDictionaryWriter m_timestamp_dict; - std::shared_ptr m_metadata_db; int m_compression_level{}; bool m_print_archive_stats{}; bool m_single_file_archive{}; diff --git a/components/core/src/clp_s/JsonParser.cpp b/components/core/src/clp_s/JsonParser.cpp index 11c615507c..ac56dc0227 100644 --- a/components/core/src/clp_s/JsonParser.cpp +++ b/components/core/src/clp_s/JsonParser.cpp @@ -127,7 +127,7 @@ JsonParser::JsonParser(JsonParserOption const& option) m_archive_options.authoritative_timestamp = m_timestamp_column; m_archive_options.authoritative_timestamp_namespace = m_timestamp_namespace; - m_archive_writer = std::make_unique(option.metadata_db); + m_archive_writer = std::make_unique(); m_archive_writer->open(m_archive_options); } diff --git a/components/job-orchestration/job_orchestration/executor/compress/compression_task.py b/components/job-orchestration/job_orchestration/executor/compress/compression_task.py index 2972c2613a..2360c148c2 100644 --- a/components/job-orchestration/job_orchestration/executor/compress/compression_task.py +++ b/components/job-orchestration/job_orchestration/executor/compress/compression_task.py @@ -10,6 +10,7 @@ from celery.app.task import Task from celery.utils.log import get_task_logger from clp_py_utils.clp_config import ( + ARCHIVES_TABLE_NAME, COMPRESSION_JOBS_TABLE_NAME, COMPRESSION_TASKS_TABLE_NAME, Database, @@ -82,6 +83,21 @@ def update_job_metadata_and_tags(db_cursor, job_id, table_prefix, tag_ids, archi ) +def update_archive_metadata(db_cursor, archive_stats): + archive_stats_defaults = { + "begin_timestamp": 0, + "end_timestamp": 0, + "creator_id": "", + "creation_ix": 0, + } + for k, v in archive_stats_defaults.items(): + archive_stats.setdefault(k, v) + keys = ", ".join(archive_stats.keys()) + values = ", ".join(f'"{v}"' if isinstance(v, str) else str(v) for v in archive_stats.values()) + query = f"INSERT INTO clp_archives ({keys}) VALUES ({values})" + db_cursor.execute(query) + + def _generate_fs_logs_list( output_file_path: pathlib.Path, paths_to_compress: PathsToCompress, @@ -347,6 +363,7 @@ def run_clp( with closing(sql_adapter.create_connection(True)) as db_conn, closing( db_conn.cursor(dictionary=True) ) as db_cursor: + update_archive_metadata(db_cursor, last_archive_stats) update_job_metadata_and_tags( db_cursor, job_id, From 106c53f3b333f983e41bf039d82d015bdc4cacd4 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Fri, 11 Apr 2025 18:18:19 -0400 Subject: [PATCH 02/14] Prevent SQL injection --- .../job_orchestration/executor/compress/compression_task.py | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/components/job-orchestration/job_orchestration/executor/compress/compression_task.py b/components/job-orchestration/job_orchestration/executor/compress/compression_task.py index 2360c148c2..edbabc4ef1 100644 --- a/components/job-orchestration/job_orchestration/executor/compress/compression_task.py +++ b/components/job-orchestration/job_orchestration/executor/compress/compression_task.py @@ -93,9 +93,9 @@ def update_archive_metadata(db_cursor, archive_stats): for k, v in archive_stats_defaults.items(): archive_stats.setdefault(k, v) keys = ", ".join(archive_stats.keys()) - values = ", ".join(f'"{v}"' if isinstance(v, str) else str(v) for v in archive_stats.values()) - query = f"INSERT INTO clp_archives ({keys}) VALUES ({values})" - db_cursor.execute(query) + value_placeholders = ", ".join(["%s"] * len(archive_stats)) + query = f"INSERT INTO {ARCHIVES_TABLE_NAME} ({keys}) VALUES ({value_placeholders})" + db_cursor.execute(query, list(archive_stats.values())) def _generate_fs_logs_list( From 985c7c26122db0c4397105c5ffbd675cc2606e38 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Fri, 11 Apr 2025 18:18:49 -0400 Subject: [PATCH 03/14] Restrict metadata update to clp-s --- .../job_orchestration/executor/compress/compression_task.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/components/job-orchestration/job_orchestration/executor/compress/compression_task.py b/components/job-orchestration/job_orchestration/executor/compress/compression_task.py index edbabc4ef1..e21c357327 100644 --- a/components/job-orchestration/job_orchestration/executor/compress/compression_task.py +++ b/components/job-orchestration/job_orchestration/executor/compress/compression_task.py @@ -363,7 +363,8 @@ def run_clp( with closing(sql_adapter.create_connection(True)) as db_conn, closing( db_conn.cursor(dictionary=True) ) as db_cursor: - update_archive_metadata(db_cursor, last_archive_stats) + if StorageEngine.CLP_S == clp_storage_engine: + update_archive_metadata(db_cursor, last_archive_stats) update_job_metadata_and_tags( db_cursor, job_id, From 1991bf4052281a3fcebf1ab1dd12255586116f3b Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Fri, 11 Apr 2025 18:30:40 -0400 Subject: [PATCH 04/14] Remove SQL database CLI from clp-s --- .../core/src/clp_s/CommandLineArguments.cpp | 29 ------------------- .../core/src/clp_s/CommandLineArguments.hpp | 8 ----- components/core/src/clp_s/clp-s.cpp | 14 --------- .../executor/compress/compression_task.py | 1 - 4 files changed, 52 deletions(-) diff --git a/components/core/src/clp_s/CommandLineArguments.cpp b/components/core/src/clp_s/CommandLineArguments.cpp index 5d3c164d86..7d19d566af 100644 --- a/components/core/src/clp_s/CommandLineArguments.cpp +++ b/components/core/src/clp_s/CommandLineArguments.cpp @@ -205,7 +205,6 @@ CommandLineArguments::parse_arguments(int argc, char const** argv) { // clang-format on po::options_description compression_options("Compression options"); - std::string metadata_db_config_file_path; std::string input_path_list_file_path; constexpr std::string_view cJsonFileType{"json"}; constexpr std::string_view cKeyValueIrFileType{"kv-ir"}; @@ -238,11 +237,6 @@ CommandLineArguments::parse_arguments(int argc, char const** argv) { po::value(&m_timestamp_key)->value_name("TIMESTAMP_COLUMN_KEY")-> default_value(m_timestamp_key), "Path (e.g. x.y) for the field containing the log event's timestamp." - )( - "db-config-file", - po::value(&metadata_db_config_file_path)->value_name("FILE")-> - default_value(metadata_db_config_file_path), - "Global metadata DB YAML config" )( "files-from,f", po::value(&input_path_list_file_path) @@ -353,29 +347,6 @@ CommandLineArguments::parse_arguments(int argc, char const** argv) { } validate_network_auth(auth, m_network_auth); - - // Parse and validate global metadata DB config - if (false == metadata_db_config_file_path.empty()) { - clp::GlobalMetadataDBConfig metadata_db_config; - try { - metadata_db_config.parse_config_file(metadata_db_config_file_path); - } catch (std::exception& e) { - SPDLOG_ERROR("Failed to validate metadata database config - {}.", e.what()); - return ParsingResult::Failure; - } - - if (clp::GlobalMetadataDBConfig::MetadataDBType::MySQL - != metadata_db_config.get_metadata_db_type()) - { - SPDLOG_ERROR( - "Invalid metadata database type for {}; only supported type is MySQL.", - m_program_name - ); - return ParsingResult::Failure; - } - - m_metadata_db_config = std::move(metadata_db_config); - } } else if ((char)Command::Extract == command_input) { po::options_description extraction_options; std::string archive_path; diff --git a/components/core/src/clp_s/CommandLineArguments.hpp b/components/core/src/clp_s/CommandLineArguments.hpp index c4a5e0dac8..a5c8295098 100644 --- a/components/core/src/clp_s/CommandLineArguments.hpp +++ b/components/core/src/clp_s/CommandLineArguments.hpp @@ -9,7 +9,6 @@ #include #include -#include "../clp/GlobalMetadataDBConfig.hpp" #include "../reducer/types.hpp" #include "Defs.hpp" #include "InputConfig.hpp" @@ -94,10 +93,6 @@ class CommandLineArguments { bool get_ignore_case() const { return m_ignore_case; } - std::optional const& get_metadata_db_config() const { - return m_metadata_db_config; - } - std::string const& get_reducer_host() const { return m_reducer_host; } int get_reducer_port() const { return m_reducer_port; } @@ -200,9 +195,6 @@ class CommandLineArguments { bool m_disable_log_order{false}; FileType m_file_type{FileType::Json}; - // Metadata db variables - std::optional m_metadata_db_config; - // MongoDB configuration variables std::string m_mongodb_uri; std::string m_mongodb_collection; diff --git a/components/core/src/clp_s/clp-s.cpp b/components/core/src/clp_s/clp-s.cpp index 360d99e9f0..e5f63013d4 100644 --- a/components/core/src/clp_s/clp-s.cpp +++ b/components/core/src/clp_s/clp-s.cpp @@ -13,7 +13,6 @@ #include #include "../clp/CurlGlobalInstance.hpp" -#include "../clp/GlobalMySQLMetadataDB.hpp" #include "../clp/streaming_archive/ArchiveMetadata.hpp" #include "../reducer/network_utils.hpp" #include "CommandLineArguments.hpp" @@ -102,19 +101,6 @@ bool compress(CommandLineArguments const& command_line_arguments) { option.structurize_arrays = command_line_arguments.get_structurize_arrays(); option.record_log_order = command_line_arguments.get_record_log_order(); - auto const& db_config_container = command_line_arguments.get_metadata_db_config(); - if (db_config_container.has_value()) { - auto const& db_config = db_config_container.value(); - option.metadata_db = std::make_shared( - db_config.get_metadata_db_host(), - db_config.get_metadata_db_port(), - db_config.get_metadata_db_username(), - db_config.get_metadata_db_password(), - db_config.get_metadata_db_name(), - db_config.get_metadata_table_prefix() - ); - } - clp_s::JsonParser parser(option); if (CommandLineArguments::FileType::KeyValueIr == option.input_file_type) { if (false == parser.parse_from_ir()) { diff --git a/components/job-orchestration/job_orchestration/executor/compress/compression_task.py b/components/job-orchestration/job_orchestration/executor/compress/compression_task.py index e21c357327..551391b8e7 100644 --- a/components/job-orchestration/job_orchestration/executor/compress/compression_task.py +++ b/components/job-orchestration/job_orchestration/executor/compress/compression_task.py @@ -198,7 +198,6 @@ def make_clp_s_command_and_env( "--target-encoded-size", str(clp_config.output.target_segment_size + clp_config.output.target_dictionaries_size), "--compression-level", str(clp_config.output.compression_level), - "--db-config-file", str(db_config_file_path), ] # fmt: on From 27a9e385a3c978f063e25512231e72581b2708ed Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Fri, 11 Apr 2025 18:40:54 -0400 Subject: [PATCH 05/14] Add comment for status update fix --- .../clp_package_utils/scripts/native/compress.py | 1 + 1 file changed, 1 insertion(+) diff --git a/components/clp-package-utils/clp_package_utils/scripts/native/compress.py b/components/clp-package-utils/clp_package_utils/scripts/native/compress.py index 2af96be38e..c5fc5720e1 100755 --- a/components/clp-package-utils/clp_package_utils/scripts/native/compress.py +++ b/components/clp-package-utils/clp_package_utils/scripts/native/compress.py @@ -83,6 +83,7 @@ def handle_job_update(db, db_cursor, job_id, no_progress_reporting): if job_last_uncompressed_size < job_uncompressed_size: print_compression_job_status(job_row, current_time) else: + # No progress in the final iteration; repeat the last status update verbatim print_compression_job_status(job_row, last_current_time) break # Done if CompressionJobStatus.FAILED == job_status: From 74c1631fa5c60890d16d5ba078440faf2198b979 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Fri, 11 Apr 2025 18:51:28 -0400 Subject: [PATCH 06/14] Remove clp::GlobalMySQLMetadataDB from JsonParser header --- components/core/src/clp_s/JsonParser.hpp | 2 -- 1 file changed, 2 deletions(-) diff --git a/components/core/src/clp_s/JsonParser.hpp b/components/core/src/clp_s/JsonParser.hpp index 4d94a9887a..aff26ff93b 100644 --- a/components/core/src/clp_s/JsonParser.hpp +++ b/components/core/src/clp_s/JsonParser.hpp @@ -17,7 +17,6 @@ #include "../clp/ffi/KeyValuePairLogEvent.hpp" #include "../clp/ffi/SchemaTree.hpp" #include "../clp/ffi/Value.hpp" -#include "../clp/GlobalMySQLMetadataDB.hpp" #include "../clp/ReaderInterface.hpp" #include "ArchiveWriter.hpp" #include "CommandLineArguments.hpp" @@ -51,7 +50,6 @@ struct JsonParserOption { bool structurize_arrays{}; bool record_log_order{true}; bool single_file_archive{false}; - std::shared_ptr metadata_db; NetworkAuthOption network_auth{}; }; From 80b2956acb7e68feb0f81a0eac8c39bc576a77d7 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Fri, 11 Apr 2025 18:56:53 -0400 Subject: [PATCH 07/14] Remove unused files from cmakelists --- components/core/src/clp_s/CMakeLists.txt | 2 -- 1 file changed, 2 deletions(-) diff --git a/components/core/src/clp_s/CMakeLists.txt b/components/core/src/clp_s/CMakeLists.txt index 2571878dfc..436794cb95 100644 --- a/components/core/src/clp_s/CMakeLists.txt +++ b/components/core/src/clp_s/CMakeLists.txt @@ -45,8 +45,6 @@ set( ../clp/GlobalMetadataDB.hpp ../clp/GlobalMetadataDBConfig.cpp ../clp/GlobalMetadataDBConfig.hpp - ../clp/GlobalMySQLMetadataDB.cpp - ../clp/GlobalMySQLMetadataDB.hpp ../clp/hash_utils.cpp ../clp/hash_utils.hpp ../clp/ir/EncodedTextAst.cpp From 5566a23eb9a28e48792d1c8e7ac3111ef7a6f09d Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Sun, 13 Apr 2025 17:24:53 -0400 Subject: [PATCH 08/14] Remove more unused files from cmakelist --- components/core/src/clp_s/CMakeLists.txt | 11 ----------- 1 file changed, 11 deletions(-) diff --git a/components/core/src/clp_s/CMakeLists.txt b/components/core/src/clp_s/CMakeLists.txt index 436794cb95..a1185baed3 100644 --- a/components/core/src/clp_s/CMakeLists.txt +++ b/components/core/src/clp_s/CMakeLists.txt @@ -18,8 +18,6 @@ set( ../clp/CurlStringList.hpp ../clp/cli_utils.cpp ../clp/cli_utils.hpp - ../clp/database_utils.cpp - ../clp/database_utils.hpp ../clp/Defs.h ../clp/ErrorCode.hpp ../clp/ffi/ir_stream/decoding_methods.cpp @@ -42,21 +40,12 @@ set( ../clp/FileDescriptor.hpp ../clp/FileReader.cpp ../clp/FileReader.hpp - ../clp/GlobalMetadataDB.hpp - ../clp/GlobalMetadataDBConfig.cpp - ../clp/GlobalMetadataDBConfig.hpp ../clp/hash_utils.cpp ../clp/hash_utils.hpp ../clp/ir/EncodedTextAst.cpp ../clp/ir/EncodedTextAst.hpp ../clp/ir/parsing.cpp ../clp/ir/parsing.hpp - ../clp/MySQLDB.cpp - ../clp/MySQLDB.hpp - ../clp/MySQLParamBindings.cpp - ../clp/MySQLParamBindings.hpp - ../clp/MySQLPreparedStatement.cpp - ../clp/MySQLPreparedStatement.hpp ../clp/NetworkReader.cpp ../clp/NetworkReader.hpp ../clp/networking/socket_utils.cpp From 91e27da4cc8f556f746921888ca9e2f7c7db36e5 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Sun, 13 Apr 2025 18:47:49 -0400 Subject: [PATCH 09/14] Remove unrelated changes --- .../clp_package_utils/scripts/native/compress.py | 8 +------- 1 file changed, 1 insertion(+), 7 deletions(-) diff --git a/components/clp-package-utils/clp_package_utils/scripts/native/compress.py b/components/clp-package-utils/clp_package_utils/scripts/native/compress.py index c5fc5720e1..b71907eb26 100755 --- a/components/clp-package-utils/clp_package_utils/scripts/native/compress.py +++ b/components/clp-package-utils/clp_package_utils/scripts/native/compress.py @@ -61,7 +61,6 @@ def handle_job_update(db, db_cursor, job_id, no_progress_reporting): ) job_last_uncompressed_size = 0 - last_current_time = datetime.datetime.now() while True: db_cursor.execute(polling_query) @@ -80,11 +79,7 @@ def handle_job_update(db, db_cursor, job_id, no_progress_reporting): # All tasks in the job is done if not no_progress_reporting: logger.info("Compression finished.") - if job_last_uncompressed_size < job_uncompressed_size: - print_compression_job_status(job_row, current_time) - else: - # No progress in the final iteration; repeat the last status update verbatim - print_compression_job_status(job_row, last_current_time) + print_compression_job_status(job_row, current_time) break # Done if CompressionJobStatus.FAILED == job_status: # One or more tasks in the job has failed @@ -103,7 +98,6 @@ def handle_job_update(db, db_cursor, job_id, no_progress_reporting): error_msg = f"Unhandled CompressionJobStatus: {job_status}" raise NotImplementedError(error_msg) - last_current_time = current_time time.sleep(0.5) From 411a40ea7bd21e959d11d362905b764af086a84b Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Sun, 13 Apr 2025 18:51:41 -0400 Subject: [PATCH 10/14] Remove unused param --- .../job_orchestration/executor/compress/compression_task.py | 2 -- 1 file changed, 2 deletions(-) diff --git a/components/job-orchestration/job_orchestration/executor/compress/compression_task.py b/components/job-orchestration/job_orchestration/executor/compress/compression_task.py index 551391b8e7..a63a4b02c3 100644 --- a/components/job-orchestration/job_orchestration/executor/compress/compression_task.py +++ b/components/job-orchestration/job_orchestration/executor/compress/compression_task.py @@ -177,7 +177,6 @@ def make_clp_s_command_and_env( clp_home: pathlib.Path, archive_output_dir: pathlib.Path, clp_config: ClpIoConfig, - db_config_file_path: pathlib.Path, use_single_file_archive: bool, ) -> Tuple[List[str], Optional[Dict[str, str]]]: """ @@ -185,7 +184,6 @@ def make_clp_s_command_and_env( :param clp_home: :param archive_output_dir: :param clp_config: - :param db_config_file_path: :param use_single_file_archive: :return: Tuple of (compression_command, compression_env_vars) """ From f3436a824ebf6bd781455eb08a3c1f3d7c1d6888 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Sun, 13 Apr 2025 19:00:11 -0400 Subject: [PATCH 11/14] Use table prefix to construct table name --- components/clp-py-utils/clp_py_utils/clp_config.py | 2 +- .../executor/compress/compression_task.py | 13 ++++++++----- 2 files changed, 9 insertions(+), 6 deletions(-) diff --git a/components/clp-py-utils/clp_py_utils/clp_config.py b/components/clp-py-utils/clp_py_utils/clp_config.py index b00c9cf673..29bf328ea5 100644 --- a/components/clp-py-utils/clp_py_utils/clp_config.py +++ b/components/clp-py-utils/clp_py_utils/clp_config.py @@ -36,7 +36,7 @@ QUERY_TASKS_TABLE_NAME = "query_tasks" COMPRESSION_JOBS_TABLE_NAME = "compression_jobs" COMPRESSION_TASKS_TABLE_NAME = "compression_tasks" -ARCHIVES_TABLE_NAME = "clp_archives" +ARCHIVES_TABLE_SUFFIX = "archives" OS_RELEASE_FILE_PATH = pathlib.Path("etc") / "os-release" diff --git a/components/job-orchestration/job_orchestration/executor/compress/compression_task.py b/components/job-orchestration/job_orchestration/executor/compress/compression_task.py index a63a4b02c3..1423672b45 100644 --- a/components/job-orchestration/job_orchestration/executor/compress/compression_task.py +++ b/components/job-orchestration/job_orchestration/executor/compress/compression_task.py @@ -10,7 +10,7 @@ from celery.app.task import Task from celery.utils.log import get_task_logger from clp_py_utils.clp_config import ( - ARCHIVES_TABLE_NAME, + ARCHIVES_TABLE_SUFFIX, COMPRESSION_JOBS_TABLE_NAME, COMPRESSION_TASKS_TABLE_NAME, Database, @@ -83,7 +83,7 @@ def update_job_metadata_and_tags(db_cursor, job_id, table_prefix, tag_ids, archi ) -def update_archive_metadata(db_cursor, archive_stats): +def update_archive_metadata(db_cursor, table_prefix, archive_stats): archive_stats_defaults = { "begin_timestamp": 0, "end_timestamp": 0, @@ -94,7 +94,9 @@ def update_archive_metadata(db_cursor, archive_stats): archive_stats.setdefault(k, v) keys = ", ".join(archive_stats.keys()) value_placeholders = ", ".join(["%s"] * len(archive_stats)) - query = f"INSERT INTO {ARCHIVES_TABLE_NAME} ({keys}) VALUES ({value_placeholders})" + query = ( + f"INSERT INTO {table_prefix}{ARCHIVES_TABLE_SUFFIX} ({keys}) VALUES ({value_placeholders})" + ) db_cursor.execute(query, list(archive_stats.values())) @@ -360,12 +362,13 @@ def run_clp( with closing(sql_adapter.create_connection(True)) as db_conn, closing( db_conn.cursor(dictionary=True) ) as db_cursor: + table_prefix = clp_metadata_db_connection_config["table_prefix"] if StorageEngine.CLP_S == clp_storage_engine: - update_archive_metadata(db_cursor, last_archive_stats) + update_archive_metadata(db_cursor, table_prefix, last_archive_stats) update_job_metadata_and_tags( db_cursor, job_id, - clp_metadata_db_connection_config["table_prefix"], + table_prefix, tag_ids, last_archive_stats, ) From 045257741f8cc4ab8734bacd3056833ec5c18240 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Sun, 13 Apr 2025 19:51:18 -0400 Subject: [PATCH 12/14] Remove unused param --- .../job_orchestration/executor/compress/compression_task.py | 1 - 1 file changed, 1 deletion(-) diff --git a/components/job-orchestration/job_orchestration/executor/compress/compression_task.py b/components/job-orchestration/job_orchestration/executor/compress/compression_task.py index 1423672b45..474f333ec1 100644 --- a/components/job-orchestration/job_orchestration/executor/compress/compression_task.py +++ b/components/job-orchestration/job_orchestration/executor/compress/compression_task.py @@ -286,7 +286,6 @@ def run_clp( clp_home=clp_home, archive_output_dir=archive_output_dir, clp_config=clp_config, - db_config_file_path=db_config_file_path, use_single_file_archive=enable_s3_write, ) else: From 03d2b51ff4ddb19e92cfc15fa47dff0ff559a6d8 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Tue, 15 Apr 2025 09:01:31 -0400 Subject: [PATCH 13/14] Address review concern --- components/core/src/clp_s/ArchiveWriter.cpp | 16 +++++++++------- components/core/src/clp_s/ArchiveWriter.hpp | 2 +- components/core/src/clp_s/CMakeLists.txt | 1 + 3 files changed, 11 insertions(+), 8 deletions(-) diff --git a/components/core/src/clp_s/ArchiveWriter.cpp b/components/core/src/clp_s/ArchiveWriter.cpp index a5e004b1cd..4e2a84882d 100644 --- a/components/core/src/clp_s/ArchiveWriter.cpp +++ b/components/core/src/clp_s/ArchiveWriter.cpp @@ -6,6 +6,7 @@ #include +#include "../clp/streaming_archive/Constants.hpp" #include "archive_constants.hpp" #include "Defs.hpp" #include "SchemaTree.hpp" @@ -426,13 +427,14 @@ std::pair ArchiveWriter::store_tables() { return {table_metadata_compressed_size, table_compressed_size}; } -void ArchiveWriter::print_archive_stats() { - nlohmann::json json_msg; - json_msg["id"] = m_id; - json_msg["begin_timestamp"] = m_timestamp_dict.get_begin_timestamp(); - json_msg["end_timestamp"] = m_timestamp_dict.get_end_timestamp(); - json_msg["uncompressed_size"] = m_uncompressed_size; - json_msg["size"] = m_compressed_size; +auto ArchiveWriter::print_archive_stats() const -> void { + using clp::streaming_archive::cMetadataDB::Archive; + nlohmann::json json_msg + = {{Archive::Id, m_id}, + {Archive::BeginTimestamp, m_timestamp_dict.get_begin_timestamp()}, + {Archive::EndTimestamp, m_timestamp_dict.get_end_timestamp()}, + {Archive::UncompressedSize, m_uncompressed_size}, + {Archive::Size, m_compressed_size}}; std::cout << json_msg.dump(-1, ' ', true, nlohmann::json::error_handler_t::ignore) << std::endl; } } // namespace clp_s diff --git a/components/core/src/clp_s/ArchiveWriter.hpp b/components/core/src/clp_s/ArchiveWriter.hpp index aa05a138fa..cd09c86df7 100644 --- a/components/core/src/clp_s/ArchiveWriter.hpp +++ b/components/core/src/clp_s/ArchiveWriter.hpp @@ -212,7 +212,7 @@ class ArchiveWriter { /** * Prints the archive's statistics (id, uncompressed size, compressed size, etc.) */ - void print_archive_stats(); + auto print_archive_stats() const -> void; static constexpr size_t cReadBlockSize = 4 * 1024; diff --git a/components/core/src/clp_s/CMakeLists.txt b/components/core/src/clp_s/CMakeLists.txt index a1185baed3..0105f2b678 100644 --- a/components/core/src/clp_s/CMakeLists.txt +++ b/components/core/src/clp_s/CMakeLists.txt @@ -57,6 +57,7 @@ set( ../clp/spdlog_with_specializations.hpp ../clp/streaming_archive/ArchiveMetadata.cpp ../clp/streaming_archive/ArchiveMetadata.hpp + ../clp/streaming_archive/Constants.hpp ../clp/streaming_compression/zstd/Decompressor.cpp ../clp/streaming_compression/zstd/Decompressor.hpp ../clp/Thread.cpp From a951bda0b1cecc3135a1286e4d3ba0fd64fb6dd4 Mon Sep 17 00:00:00 2001 From: Bingran Hu Date: Tue, 15 Apr 2025 10:05:25 -0400 Subject: [PATCH 14/14] Fix namespace usage --- components/core/src/clp_s/ArchiveWriter.cpp | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/components/core/src/clp_s/ArchiveWriter.cpp b/components/core/src/clp_s/ArchiveWriter.cpp index 4e2a84882d..a05856da16 100644 --- a/components/core/src/clp_s/ArchiveWriter.cpp +++ b/components/core/src/clp_s/ArchiveWriter.cpp @@ -428,7 +428,7 @@ std::pair ArchiveWriter::store_tables() { } auto ArchiveWriter::print_archive_stats() const -> void { - using clp::streaming_archive::cMetadataDB::Archive; + namespace Archive = clp::streaming_archive::cMetadataDB::Archive; nlohmann::json json_msg = {{Archive::Id, m_id}, {Archive::BeginTimestamp, m_timestamp_dict.get_begin_timestamp()},