MCPcopy Create free account
hub / github.com/apache/impala / CreateAndOpenScannerHelper

Method CreateAndOpenScannerHelper

be/src/exec/hdfs-scan-node-base.cc:909–984  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

907}
908
909Status 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_));

Callers

nothing calls this directly

Calls 10

SubstituteFunction · 0.85
file_descMethod · 0.80
GetStreamMethod · 0.80
StatusClass · 0.70
getMethod · 0.65
resetMethod · 0.65
file_formatMethod · 0.45
findMethod · 0.45
endMethod · 0.45
OpenMethod · 0.45

Tested by

no test coverage detected