| 907 | } |
| 908 | |
| 909 | Status HdfsScanNodeBase::CreateAndOpenScannerHelper(HdfsPartitionDescriptor* partition, |
| 910 | ScannerContext* context, scoped_ptr<HdfsScanner>* scanner) { |
| 911 | using namespace org::apache::impala::fb; |
| 912 | DCHECK(context != nullptr); |
| 913 | DCHECK(scanner->get() == nullptr); |
| 914 | |
| 915 | const FbSplitFileMetadata* file_metadata = |
| 916 | context->GetStream(0)->file_desc()->file_metadata; |
| 917 | if (file_metadata) { |
| 918 | // Iceberg tables can have different file format for each data file: |
| 919 | const FbIcebergSplitMetadata* ice_metadata = file_metadata->iceberg_metadata(); |
| 920 | DCHECK(ice_metadata != nullptr); |
| 921 | switch (ice_metadata->file_format()) { |
| 922 | case FbIcebergDataFileFormat::FbIcebergDataFileFormat_PARQUET: |
| 923 | scanner->reset(new HdfsParquetScanner(this, runtime_state_)); |
| 924 | break; |
| 925 | case FbIcebergDataFileFormat::FbIcebergDataFileFormat_ORC: |
| 926 | scanner->reset(new HdfsOrcScanner(this, runtime_state_)); |
| 927 | break; |
| 928 | case FbIcebergDataFileFormat::FbIcebergDataFileFormat_AVRO: |
| 929 | scanner->reset(new HdfsAvroScanner(this, runtime_state_)); |
| 930 | break; |
| 931 | default: |
| 932 | return Status(Substitute( |
| 933 | "Unknown Iceberg file format type: $0", ice_metadata->file_format())); |
| 934 | } |
| 935 | } else { |
| 936 | THdfsCompression::type compression = |
| 937 | context->GetStream()->file_desc()->file_compression; |
| 938 | |
| 939 | // Create a new scanner for this file format and compression. |
| 940 | switch (partition->file_format()) { |
| 941 | case THdfsFileFormat::TEXT: |
| 942 | if (HdfsTextScanner::HasBuiltinSupport(compression)) { |
| 943 | scanner->reset(new HdfsTextScanner(this, runtime_state_)); |
| 944 | } else { |
| 945 | // No builtin support - we must have loaded the plugin in IssueInitialRanges(). |
| 946 | auto it = _THdfsCompression_VALUES_TO_NAMES.find(compression); |
| 947 | DCHECK(it != _THdfsCompression_VALUES_TO_NAMES.end()) |
| 948 | << "Already issued ranges for this compression type."; |
| 949 | scanner->reset(HdfsPluginTextScanner::GetHdfsPluginTextScanner( |
| 950 | this, runtime_state_, it->second)); |
| 951 | } |
| 952 | break; |
| 953 | case THdfsFileFormat::SEQUENCE_FILE: |
| 954 | scanner->reset(new HdfsSequenceScanner(this, runtime_state_)); |
| 955 | break; |
| 956 | case THdfsFileFormat::RC_FILE: |
| 957 | scanner->reset(new HdfsRCFileScanner(this, runtime_state_)); |
| 958 | break; |
| 959 | case THdfsFileFormat::AVRO: |
| 960 | scanner->reset(new HdfsAvroScanner(this, runtime_state_)); |
| 961 | break; |
| 962 | case THdfsFileFormat::PARQUET: |
| 963 | scanner->reset(new HdfsParquetScanner(this, runtime_state_)); |
| 964 | break; |
| 965 | case THdfsFileFormat::ORC: |
| 966 | scanner->reset(new HdfsOrcScanner(this, runtime_state_)); |
nothing calls this directly
no test coverage detected