| 73 | } |
| 74 | |
| 75 | Status DataSinkConfig::CreateConfig(const TDataSink& thrift_sink, |
| 76 | const RowDescriptor* row_desc, FragmentState* state, |
| 77 | DataSinkConfig** data_sink) { |
| 78 | ObjectPool* pool = state->obj_pool(); |
| 79 | *data_sink = nullptr; |
| 80 | switch (thrift_sink.type) { |
| 81 | case TDataSinkType::DATA_STREAM_SINK: |
| 82 | if (!thrift_sink.__isset.stream_sink) return Status("Missing data stream sink."); |
| 83 | // TODO: figure out good buffer size based on size of output row |
| 84 | *data_sink = pool->Add(new KrpcDataStreamSenderConfig()); |
| 85 | break; |
| 86 | case TDataSinkType::TABLE_SINK: |
| 87 | if (!thrift_sink.__isset.table_sink) return Status("Missing table sink."); |
| 88 | switch (thrift_sink.table_sink.type) { |
| 89 | case TTableSinkType::HDFS: |
| 90 | if (thrift_sink.table_sink.action == TSinkAction::INSERT) { |
| 91 | *data_sink = pool->Add(new HdfsTableSinkConfig()); |
| 92 | } else if (thrift_sink.table_sink.action == TSinkAction::DELETE) { |
| 93 | // Currently only Iceberg tables support DELETE operations for FS tables. |
| 94 | *data_sink = pool->Add(new IcebergDeleteSinkConfig()); |
| 95 | } |
| 96 | break; |
| 97 | case TTableSinkType::KUDU: |
| 98 | RETURN_IF_ERROR(CheckKuduAvailability()); |
| 99 | *data_sink = pool->Add(new KuduTableSinkConfig()); |
| 100 | break; |
| 101 | case TTableSinkType::HBASE: |
| 102 | *data_sink = pool->Add(new HBaseTableSinkConfig()); |
| 103 | break; |
| 104 | default: |
| 105 | stringstream error_msg; |
| 106 | map<int, const char*>::const_iterator i = |
| 107 | _TTableSinkType_VALUES_TO_NAMES.find(thrift_sink.table_sink.type); |
| 108 | const char* str = i != _TTableSinkType_VALUES_TO_NAMES.end() ? |
| 109 | i->second : |
| 110 | "Unknown table sink"; |
| 111 | error_msg << str << " not implemented."; |
| 112 | return Status(error_msg.str()); |
| 113 | } |
| 114 | break; |
| 115 | case TDataSinkType::PLAN_ROOT_SINK: |
| 116 | *data_sink = pool->Add(new PlanRootSinkConfig()); |
| 117 | break; |
| 118 | case TDataSinkType::HASH_JOIN_BUILDER: { |
| 119 | *data_sink = pool->Add(new PhjBuilderConfig()); |
| 120 | break; |
| 121 | } |
| 122 | case TDataSinkType::NESTED_LOOP_JOIN_BUILDER: { |
| 123 | *data_sink = pool->Add(new NljBuilderConfig()); |
| 124 | break; |
| 125 | } |
| 126 | case TDataSinkType::ICEBERG_DELETE_BUILDER: { |
| 127 | *data_sink = pool->Add(new IcebergDeleteBuilderConfig()); |
| 128 | break; |
| 129 | } |
| 130 | case TDataSinkType::MULTI_DATA_SINK: { |
| 131 | *data_sink = pool->Add(new MultiTableSinkConfig()); |
| 132 | break; |