| 2873 | } |
| 2874 | |
| 2875 | Result<ExecBatch> KeyPayloadFromInput(int side, ExecBatch* input) { |
| 2876 | ExecBatch projected({}, input->length); |
| 2877 | int num_key_cols = schema_[side]->num_cols(HashJoinProjection::KEY); |
| 2878 | int num_payload_cols = schema_[side]->num_cols(HashJoinProjection::PAYLOAD); |
| 2879 | projected.values.resize(num_key_cols + num_payload_cols); |
| 2880 | |
| 2881 | auto key_to_input = |
| 2882 | schema_[side]->map(HashJoinProjection::KEY, HashJoinProjection::INPUT); |
| 2883 | for (int icol = 0; icol < num_key_cols; ++icol) { |
| 2884 | const Datum& value_in = input->values[key_to_input.get(icol)]; |
| 2885 | if (value_in.is_scalar()) { |
| 2886 | ARROW_ASSIGN_OR_RAISE( |
| 2887 | projected.values[icol], |
| 2888 | MakeArrayFromScalar(*value_in.scalar(), projected.length, pool_)); |
| 2889 | } else { |
| 2890 | projected.values[icol] = value_in; |
| 2891 | } |
| 2892 | } |
| 2893 | auto payload_to_input = |
| 2894 | schema_[side]->map(HashJoinProjection::PAYLOAD, HashJoinProjection::INPUT); |
| 2895 | for (int icol = 0; icol < num_payload_cols; ++icol) { |
| 2896 | const Datum& value_in = input->values[payload_to_input.get(icol)]; |
| 2897 | if (value_in.is_scalar()) { |
| 2898 | ARROW_ASSIGN_OR_RAISE( |
| 2899 | projected.values[num_key_cols + icol], |
| 2900 | MakeArrayFromScalar(*value_in.scalar(), projected.length, pool_)); |
| 2901 | } else { |
| 2902 | projected.values[num_key_cols + icol] = value_in; |
| 2903 | } |
| 2904 | } |
| 2905 | |
| 2906 | return projected; |
| 2907 | } |
| 2908 | |
| 2909 | bool IsCancelled() { return cancelled_.load(); } |
| 2910 |
nothing calls this directly
no test coverage detected