| 836 | |
| 837 | #[tokio::test] |
| 838 | async fn test_dictionary_hydration() { |
| 839 | let arr1: DictionaryArray<UInt16Type> = vec!["a", "a", "b"].into_iter().collect(); |
| 840 | let arr2: DictionaryArray<UInt16Type> = vec!["c", "c", "d"].into_iter().collect(); |
| 841 | |
| 842 | let schema = Arc::new(Schema::new(vec![Field::new_dictionary( |
| 843 | "dict", |
| 844 | DataType::UInt16, |
| 845 | DataType::Utf8, |
| 846 | false, |
| 847 | )])); |
| 848 | let batch1 = RecordBatch::try_new(schema.clone(), vec![Arc::new(arr1)]).unwrap(); |
| 849 | let batch2 = RecordBatch::try_new(schema, vec![Arc::new(arr2)]).unwrap(); |
| 850 | |
| 851 | let stream = futures::stream::iter(vec![Ok(batch1), Ok(batch2)]); |
| 852 | |
| 853 | let encoder = FlightDataEncoderBuilder::default().build(stream); |
| 854 | let mut decoder = FlightDataDecoder::new(encoder); |
| 855 | let expected_schema = Schema::new(vec![Field::new("dict", DataType::Utf8, false)]); |
| 856 | let expected_schema = Arc::new(expected_schema); |
| 857 | let mut expected_arrays = vec![ |
| 858 | StringArray::from(vec!["a", "a", "b"]), |
| 859 | StringArray::from(vec!["c", "c", "d"]), |
| 860 | ] |
| 861 | .into_iter(); |
| 862 | while let Some(decoded) = decoder.next().await { |
| 863 | let decoded = decoded.unwrap(); |
| 864 | match decoded.payload { |
| 865 | DecodedPayload::None => {} |
| 866 | DecodedPayload::Schema(s) => assert_eq!(s, expected_schema), |
| 867 | DecodedPayload::RecordBatch(b) => { |
| 868 | assert_eq!(b.schema(), expected_schema); |
| 869 | let expected_array = expected_arrays.next().unwrap(); |
| 870 | let actual_array = b.column_by_name("dict").unwrap(); |
| 871 | let actual_array = downcast_array::<StringArray>(actual_array); |
| 872 | |
| 873 | assert_eq!(actual_array, expected_array); |
| 874 | } |
| 875 | } |
| 876 | } |
| 877 | } |
| 878 | |
| 879 | #[tokio::test] |
| 880 | async fn test_dictionary_resend() { |