| 638 | } |
| 639 | |
| 640 | static void logQueryFinishImpl( |
| 641 | QueryLogElement & elem, |
| 642 | const ContextMutablePtr & context, |
| 643 | const ASTPtr & query_ast, |
| 644 | const QueryPipelineFinalizedInfo & query_pipeline_finalized_info, |
| 645 | bool pulling_pipeline, |
| 646 | std::shared_ptr<OpenTelemetry::SpanHolder> query_span, |
| 647 | QueryResultCacheUsage query_result_cache_usage, |
| 648 | bool internal, |
| 649 | std::chrono::system_clock::time_point time) |
| 650 | { |
| 651 | const Settings & settings = context->getSettingsRef(); |
| 652 | auto log_queries = settings[Setting::log_queries]; |
| 653 | |
| 654 | if (QueryStatusPtr process_list_elem = context->getProcessListElement()) |
| 655 | { |
| 656 | { |
| 657 | ResultProgress result_progress(0, 0, 0); |
| 658 | |
| 659 | chassert((query_pipeline_finalized_info.result_progress != std::nullopt) == pulling_pipeline); |
| 660 | |
| 661 | if (query_pipeline_finalized_info.result_progress) |
| 662 | { |
| 663 | result_progress = *query_pipeline_finalized_info.result_progress; |
| 664 | } |
| 665 | else if (!pulling_pipeline) |
| 666 | { |
| 667 | auto progress_out = process_list_elem->getProgressOut(); |
| 668 | result_progress.result_rows = progress_out.written_rows; |
| 669 | result_progress.result_bytes = progress_out.written_bytes; |
| 670 | } |
| 671 | |
| 672 | if (auto progress_callback = context->getProgressCallback()) |
| 673 | { |
| 674 | Progress p; |
| 675 | p.incrementPiecewiseAtomically(Progress{result_progress}); |
| 676 | progress_callback(p); |
| 677 | } |
| 678 | |
| 679 | elem.result_rows = result_progress.result_rows; |
| 680 | elem.result_bytes = result_progress.result_bytes; |
| 681 | } |
| 682 | |
| 683 | QueryStatusInfo info = process_list_elem->getInfo(true, settings[Setting::log_profile_events]); |
| 684 | logQueryMetricLogFinish(context, internal, elem.client_info.current_query_id, time, std::make_shared<QueryStatusInfo>(info)); |
| 685 | |
| 686 | elem.type = QueryLogElementType::QUERY_FINISH; |
| 687 | |
| 688 | addStatusInfoToQueryLogElement(elem, info, query_ast, context, time); |
| 689 | |
| 690 | if (elem.read_rows != 0) |
| 691 | { |
| 692 | double elapsed_seconds = static_cast<double>(info.elapsed_microseconds) / 1000000.0; |
| 693 | double rows_per_second = static_cast<double>(elem.read_rows) / elapsed_seconds; |
| 694 | double bytes_per_second = static_cast<double>(elem.read_bytes) / elapsed_seconds; |
| 695 | LOG_DEBUG( |
| 696 | getLogger("executeQuery"), |
| 697 | "Read {} rows, {} in {} sec., {} rows/sec., {}/sec.", |
no test coverage detected