@@ -2566,6 +2566,237 @@ TEST_P(ScanAndReadInteTest, TestCastTimestampType) {
25662566 ASSERT_TRUE (expected->Equals (read_result)) << read_result->ToString ();
25672567}
25682568
2569+ TEST_P (ScanAndReadInteTest, TestAvroWithAppendSnapshot1) {
2570+ auto [file_format, enable_prefetch] = GetParam ();
2571+ if (file_format != " avro" ) {
2572+ return ;
2573+ }
2574+ std::string table_path = GetDataDir () + " /avro/append_multiple.db/append_multiple" ;
2575+
2576+ // scan
2577+ ScanContextBuilder scan_context_builder (table_path);
2578+ scan_context_builder.AddOption (Options::SCAN_SNAPSHOT_ID , " 1" );
2579+ ASSERT_OK_AND_ASSIGN (auto scan_context, scan_context_builder.Finish ());
2580+ ASSERT_OK_AND_ASSIGN (auto table_scan, TableScan::Create (std::move (scan_context)));
2581+ ASSERT_OK_AND_ASSIGN (auto result_plan, table_scan->CreatePlan ());
2582+ ASSERT_EQ (result_plan->SnapshotId ().value (), 1 );
2583+
2584+ auto splits = result_plan->Splits ();
2585+ ASSERT_EQ (3 , splits.size ());
2586+
2587+ // read
2588+ ReadContextBuilder read_context_builder (table_path);
2589+ ASSERT_OK_AND_ASSIGN (std::unique_ptr<ReadContext> read_context, read_context_builder.Finish ());
2590+ ASSERT_OK_AND_ASSIGN (auto table_read, TableRead::Create (std::move (read_context)));
2591+ ASSERT_OK_AND_ASSIGN (auto batch_reader, table_read->CreateReader (splits));
2592+ ASSERT_OK_AND_ASSIGN (auto read_result, ReadResultCollector::CollectResult (batch_reader.get ()));
2593+
2594+ // check result
2595+ auto timezone = DateTimeUtils::GetLocalTimezoneName ();
2596+ arrow::FieldVector fields = {
2597+ arrow::field (" _VALUE_KIND" , arrow::int8 ()),
2598+ arrow::field (" f0" , arrow::int8 ()),
2599+ arrow::field (" f1" , arrow::int16 ()),
2600+ arrow::field (" f2" , arrow::int32 ()),
2601+ arrow::field (" f3" , arrow::int64 ()),
2602+ arrow::field (" f4" , arrow::float32 ()),
2603+ arrow::field (" f5" , arrow::float64 ()),
2604+ arrow::field (" f6" , arrow::utf8 ()),
2605+ arrow::field (" f7" , arrow::binary ()),
2606+ arrow::field (" f8" , arrow::date32 ()),
2607+ arrow::field (" f9" , arrow::decimal128 (5 , 2 )),
2608+ arrow::field (" f10" , arrow::timestamp (arrow::TimeUnit::SECOND )),
2609+ arrow::field (" f11" , arrow::timestamp (arrow::TimeUnit::MILLI )),
2610+ arrow::field (" f12" , arrow::timestamp (arrow::TimeUnit::MICRO )),
2611+ arrow::field (" f13" , arrow::timestamp (arrow::TimeUnit::SECOND , timezone)),
2612+ arrow::field (" f14" , arrow::timestamp (arrow::TimeUnit::MILLI , timezone)),
2613+ arrow::field (" f15" , arrow::timestamp (arrow::TimeUnit::MICRO , timezone)),
2614+ arrow::field (" f16" ,
2615+ arrow::struct_ ({arrow::field (" f0" , arrow::map (arrow::utf8 (), arrow::int32 ())),
2616+ arrow::field (" f1" , arrow::list (arrow::int32 ()))})),
2617+ };
2618+ auto expected = std::make_shared<arrow::ChunkedArray>(
2619+ arrow::ipc::internal::json::ArrayFromJSON (arrow::struct_ (fields), R"( [
2620+ [0, 2, 10, 1, 100, 2.0, 2.0, "two", "bbb", 123, "123.45", "1970-01-02 00:00:00", "1970-01-02 00:00:00.000", "1970-01-02 00:00:00.000000", "1970-01-02 00:00:00", "1970-01-02 00:00:00.000", "1970-01-02 00:00:00.000000",[[["key",123]],[1,2,3]]],
2621+ [0, 1, 10, 0, 100, 1.0, 1.0, "one", "aaa", 123, "123.45", "1970-01-01 00:00:00", "1970-01-01 00:00:00.000", "1970-01-01 00:00:00.000000", "1970-01-01 00:00:00", "1970-01-01 00:00:00.000", "1970-01-01 00:00:00.000000",[[["key",123]],[1,2,3]]],
2622+ [0, 3, 11, 0, 100, null, 3.0, "three", "ccc", 123, "123.45", "1970-01-03 00:00:00", "1970-01-03 00:00:00.000", "1970-01-03 00:00:00.000000", "1970-01-03 00:00:00", "1970-01-03 00:00:00.000", "1970-01-03 00:00:00.000000",[[["key",123]],[1,2,3]]],
2623+ [0, 4, 11, 0, 100, 4.0, null, "four", "ddd", 123, "123.45", "1970-01-04 00:00:00", "1970-01-04 00:00:00.000", "1970-01-04 00:00:00.000000", "1970-01-04 00:00:00", "1970-01-04 00:00:00.000", "1970-01-04 00:00:00.000000",[[["key",123]],[1,2,3]]]
2624+ ])" )
2625+ .ValueOrDie ());
2626+ ASSERT_TRUE (expected);
2627+ ASSERT_TRUE (expected->Equals (read_result)) << read_result->ToString ();
2628+ }
2629+
2630+ TEST_P (ScanAndReadInteTest, TestAvroWithAppendSnapshot2) {
2631+ auto [file_format, enable_prefetch] = GetParam ();
2632+ if (file_format != " avro" ) {
2633+ return ;
2634+ }
2635+ std::string table_path = GetDataDir () + " /avro/append_multiple.db/append_multiple" ;
2636+
2637+ // scan
2638+ ScanContextBuilder scan_context_builder (table_path);
2639+ scan_context_builder.AddOption (Options::SCAN_SNAPSHOT_ID , " 2" );
2640+ ASSERT_OK_AND_ASSIGN (auto scan_context, scan_context_builder.Finish ());
2641+ ASSERT_OK_AND_ASSIGN (auto table_scan, TableScan::Create (std::move (scan_context)));
2642+ ASSERT_OK_AND_ASSIGN (auto result_plan, table_scan->CreatePlan ());
2643+ ASSERT_EQ (result_plan->SnapshotId ().value (), 2 );
2644+
2645+ auto splits = result_plan->Splits ();
2646+ ASSERT_EQ (3 , splits.size ());
2647+
2648+ // read
2649+ ReadContextBuilder read_context_builder (table_path);
2650+ ASSERT_OK_AND_ASSIGN (std::unique_ptr<ReadContext> read_context, read_context_builder.Finish ());
2651+ ASSERT_OK_AND_ASSIGN (auto table_read, TableRead::Create (std::move (read_context)));
2652+ ASSERT_OK_AND_ASSIGN (auto batch_reader, table_read->CreateReader (splits));
2653+ ASSERT_OK_AND_ASSIGN (auto read_result, ReadResultCollector::CollectResult (batch_reader.get ()));
2654+
2655+ // check result
2656+ auto timezone = DateTimeUtils::GetLocalTimezoneName ();
2657+ arrow::FieldVector fields = {
2658+ arrow::field (" _VALUE_KIND" , arrow::int8 ()),
2659+ arrow::field (" f0" , arrow::int8 ()),
2660+ arrow::field (" f1" , arrow::int16 ()),
2661+ arrow::field (" f2" , arrow::int32 ()),
2662+ arrow::field (" f3" , arrow::int64 ()),
2663+ arrow::field (" f4" , arrow::float32 ()),
2664+ arrow::field (" f5" , arrow::float64 ()),
2665+ arrow::field (" f6" , arrow::utf8 ()),
2666+ arrow::field (" f7" , arrow::binary ()),
2667+ arrow::field (" f8" , arrow::date32 ()),
2668+ arrow::field (" f9" , arrow::decimal128 (5 , 2 )),
2669+ arrow::field (" f10" , arrow::timestamp (arrow::TimeUnit::SECOND )),
2670+ arrow::field (" f11" , arrow::timestamp (arrow::TimeUnit::MILLI )),
2671+ arrow::field (" f12" , arrow::timestamp (arrow::TimeUnit::MICRO )),
2672+ arrow::field (" f13" , arrow::timestamp (arrow::TimeUnit::SECOND , timezone)),
2673+ arrow::field (" f14" , arrow::timestamp (arrow::TimeUnit::MILLI , timezone)),
2674+ arrow::field (" f15" , arrow::timestamp (arrow::TimeUnit::MICRO , timezone)),
2675+ arrow::field (" f16" ,
2676+ arrow::struct_ ({arrow::field (" f0" , arrow::map (arrow::utf8 (), arrow::int32 ())),
2677+ arrow::field (" f1" , arrow::list (arrow::int32 ()))})),
2678+ };
2679+ auto expected = std::make_shared<arrow::ChunkedArray>(
2680+ arrow::ipc::internal::json::ArrayFromJSON (arrow::struct_ (fields), R"( [
2681+ [0, 6, 10, 1, 100, 6.0, 4.0, "six", "fff", 123, "123.45", "1970-01-02 00:00:00", "1970-01-06 00:00:00.000", "1970-01-06 00:00:00.000000", "1970-01-06 00:00:00", "1970-01-06 00:00:00.000", "1970-01-06 00:00:00.000000",[[["key",123]],[1,2,3]]],
2682+ [0, 5, 10, 0, 100, 5.0, 2.0, null, "eee", 123, "123.45", "1970-01-01 00:00:00", "1970-01-05 00:00:00.000", "1970-01-05 00:00:00.000000", "1970-01-05 00:00:00", "1970-01-05 00:00:00.000", "1970-01-05 00:00:00.000000",[[["key",123]],[1,2,3]]],
2683+ [0, 7, 11, 0, 100, 7.0, 6.0, "seven", "ggg", 123, "123.45", "1970-01-03 00:00:00", "1970-01-07 00:00:00.000", "1970-01-07 00:00:00.000000", "1970-01-07 00:00:00", "1970-01-07 00:00:00.000", "1970-01-07 00:00:00.000000",[[["key",123]],[1,2,3]]]
2684+ ])" )
2685+ .ValueOrDie ());
2686+ ASSERT_TRUE (expected);
2687+ ASSERT_TRUE (expected->Equals (read_result)) << read_result->ToString ();
2688+ }
2689+
2690+ TEST_P (ScanAndReadInteTest, TestAvroWithPkSnapshot1) {
2691+ auto [file_format, enable_prefetch] = GetParam ();
2692+ if (file_format != " avro" ) {
2693+ return ;
2694+ }
2695+ std::string table_path = GetDataDir () + " /avro/pk_with_multiple_type.db/pk_with_multiple_type" ;
2696+
2697+ // scan
2698+ ScanContextBuilder scan_context_builder (table_path);
2699+ scan_context_builder.AddOption (Options::SCAN_SNAPSHOT_ID , " 1" );
2700+ ASSERT_OK_AND_ASSIGN (auto scan_context, scan_context_builder.Finish ());
2701+ ASSERT_OK_AND_ASSIGN (auto table_scan, TableScan::Create (std::move (scan_context)));
2702+ ASSERT_OK_AND_ASSIGN (auto result_plan, table_scan->CreatePlan ());
2703+ ASSERT_EQ (result_plan->SnapshotId ().value (), 1 );
2704+
2705+ auto splits = result_plan->Splits ();
2706+ ASSERT_EQ (1 , splits.size ());
2707+
2708+ // read
2709+ ReadContextBuilder read_context_builder (table_path);
2710+ ASSERT_OK_AND_ASSIGN (std::unique_ptr<ReadContext> read_context, read_context_builder.Finish ());
2711+ ASSERT_OK_AND_ASSIGN (auto table_read, TableRead::Create (std::move (read_context)));
2712+ ASSERT_OK_AND_ASSIGN (auto batch_reader, table_read->CreateReader (splits));
2713+ ASSERT_OK_AND_ASSIGN (auto read_result, ReadResultCollector::CollectResult (batch_reader.get ()));
2714+
2715+ // check result
2716+ arrow::FieldVector fields = {
2717+ arrow::field (" _VALUE_KIND" , arrow::int8 ()),
2718+ arrow::field (" f0" , arrow::boolean ()),
2719+ arrow::field (" f1" , arrow::int8 ()),
2720+ arrow::field (" f2" , arrow::int16 ()),
2721+ arrow::field (" f3" , arrow::int32 ()),
2722+ arrow::field (" f4" , arrow::int64 ()),
2723+ arrow::field (" f5" , arrow::float32 ()),
2724+ arrow::field (" f6" , arrow::float64 ()),
2725+ arrow::field (" f7" , arrow::utf8 ()),
2726+ arrow::field (" f8" , arrow::binary ()),
2727+ arrow::field (" f9" , arrow::date32 ()),
2728+ arrow::field (" f10" , arrow::decimal128 (5 , 2 )),
2729+ arrow::field (" f11" ,
2730+ arrow::struct_ ({arrow::field (" f0" , arrow::map (arrow::utf8 (), arrow::int32 ())),
2731+ arrow::field (" f1" , arrow::list (arrow::int32 ()))})),
2732+ };
2733+ auto expected = std::make_shared<arrow::ChunkedArray>(
2734+ arrow::ipc::internal::json::ArrayFromJSON (struct_ (fields), R"( [
2735+ [0, false, 10, 1, 1, 1000, 1.5, 2.5, "Alice", "abcdef", 100, "123.45", [[["key",123]],[1,2,3]]],
2736+ [0, false, 10, 1, 1, 1000, 1.5, 2.5, "Bob", "abcdef", 100, "123.45", [[["key",123]],[1,2,3]]],
2737+ [0, true, 10, 1, 1, 1000, 1.5, 2.5, "Emily", "abcdef", 100, "123.45", [[["key",123]],[1,2,3]]],
2738+ [0, true, 10, 1, 1, 1000, 1.5, 2.5, "Tony", "abcdef", 100, "123.45", [[["key",123]],[1,2,3]]]
2739+ ])" )
2740+ .ValueOrDie ());
2741+ ASSERT_TRUE (expected);
2742+ ASSERT_TRUE (expected->Equals (read_result)) << read_result->ToString ();
2743+ }
2744+
2745+ TEST_P (ScanAndReadInteTest, TestAvroWithPkSnapshot2) {
2746+ auto [file_format, enable_prefetch] = GetParam ();
2747+ if (file_format != " avro" ) {
2748+ return ;
2749+ }
2750+ std::string table_path = GetDataDir () + " /avro/pk_with_multiple_type.db/pk_with_multiple_type" ;
2751+
2752+ // scan
2753+ ScanContextBuilder scan_context_builder (table_path);
2754+ scan_context_builder.AddOption (Options::SCAN_SNAPSHOT_ID , " 2" );
2755+ ASSERT_OK_AND_ASSIGN (auto scan_context, scan_context_builder.Finish ());
2756+ ASSERT_OK_AND_ASSIGN (auto table_scan, TableScan::Create (std::move (scan_context)));
2757+ ASSERT_OK_AND_ASSIGN (auto result_plan, table_scan->CreatePlan ());
2758+ ASSERT_EQ (result_plan->SnapshotId ().value (), 2 );
2759+
2760+ auto splits = result_plan->Splits ();
2761+ ASSERT_EQ (1 , splits.size ());
2762+
2763+ // read
2764+ ReadContextBuilder read_context_builder (table_path);
2765+ ASSERT_OK_AND_ASSIGN (std::unique_ptr<ReadContext> read_context, read_context_builder.Finish ());
2766+ ASSERT_OK_AND_ASSIGN (auto table_read, TableRead::Create (std::move (read_context)));
2767+ ASSERT_OK_AND_ASSIGN (auto batch_reader, table_read->CreateReader (splits));
2768+ ASSERT_OK_AND_ASSIGN (auto read_result, ReadResultCollector::CollectResult (batch_reader.get ()));
2769+
2770+ // check result
2771+ arrow::FieldVector fields = {
2772+ arrow::field (" _VALUE_KIND" , arrow::int8 ()),
2773+ arrow::field (" f0" , arrow::boolean ()),
2774+ arrow::field (" f1" , arrow::int8 ()),
2775+ arrow::field (" f2" , arrow::int16 ()),
2776+ arrow::field (" f3" , arrow::int32 ()),
2777+ arrow::field (" f4" , arrow::int64 ()),
2778+ arrow::field (" f5" , arrow::float32 ()),
2779+ arrow::field (" f6" , arrow::float64 ()),
2780+ arrow::field (" f7" , arrow::utf8 ()),
2781+ arrow::field (" f8" , arrow::binary ()),
2782+ arrow::field (" f9" , arrow::date32 ()),
2783+ arrow::field (" f10" , arrow::decimal128 (5 , 2 )),
2784+ arrow::field (" f11" ,
2785+ arrow::struct_ ({arrow::field (" f0" , arrow::map (arrow::utf8 (), arrow::int32 ())),
2786+ arrow::field (" f1" , arrow::list (arrow::int32 ()))})),
2787+ };
2788+ auto expected = std::make_shared<arrow::ChunkedArray>(
2789+ arrow::ipc::internal::json::ArrayFromJSON (struct_ (fields), R"( [
2790+ [0, false, 10, 1, 1, 1000, 1.5, 2.5, "Alice", "abcdef", 100, "123.45", [[["key",123]],[1,2,3]]],
2791+ [0, false, 10, 1, 1, 1000, 1.5, 2.5, "Bob", "abcdef", 100, "123.45", [[["key",123]],[1,2,3]]],
2792+ [0, true, 10, 1, 1, 1000, 1.5, 2.5, "Lucy", "abcdef", 100, "123.45", [[["key",123]],[1,2,3]]],
2793+ [0, true, 10, 1, 1, 1000, 1.5, 2.5, "Tony", "abcdef", 100, "123.45", [[["key",123]],[1,2,3]]]
2794+ ])" )
2795+ .ValueOrDie ());
2796+ ASSERT_TRUE (expected);
2797+ ASSERT_TRUE (expected->Equals (read_result)) << read_result->ToString ();
2798+ }
2799+
25692800std::vector<std::pair<std::string, bool >> GetTestValuesForScanAndReadInteTest () {
25702801 std::vector<std::pair<std::string, bool >> values = {{" parquet" , false }, {" parquet" , true }};
25712802#ifdef PAIMON_ENABLE_ORC
0 commit comments