vbhanuchander-lang commented on PR #17644:
URL: https://github.com/apache/iceberg/pull/17644#issuecomment-5837524661

   Thanks @pvary and @Guosmilesmile — all of it is addressed. Rebased onto 
current `main` and **restricted to `flink/v2.3`**, so this is now 5 files 
instead of 16; I will backport to v1.20, v2.1 and v2.2 once you are happy with 
the shape here.
   
   ### @pvary
   
   **Duplicated code extracted.** The multiset-to-map conversion was repeated 
in all three visitors, so it is now one package-private `FlinkMultisets` with 
`isMapLike` and `asMapType`. That also removed the special-casing inside 
`AvroWithFlinkSchemaVisitor#mapKeyType`/`mapValueType`, which are now one line 
each, and let `ParquetWithFlinkSchemaVisitor` drop its private `toMapType`. The 
production diff outside the new file is down to 19 lines.
   
   **Error messages** now say `is not a map or multiset`, in the two 
`AvroWithFlinkSchemaVisitor` preconditions and the 
`ParquetWithFlinkSchemaVisitor` one.
   
   **"All of the tests should check the data read back"** — they do now. Every 
test writes and then reads back through the matching reader and compares `id` 
and the tag counts, rather than asserting `doesNotThrowAnyException`. That was 
the right call for a second reason: 
`assertThatCode(...).doesNotThrowAnyException()` would have passed even if the 
counts had been written wrongly.
   
   **Test conventions** — class and methods are package-private, the `test` 
prefix is gone, and the javadoc describes what is true now rather than what 
used to fail. The fully-qualified `java.io.File`/`java.nio.file.Files` in the 
ORC test are proper imports.
   
   **`MULTISET<STRING>` nullable** added as 
`parquetRoundTripsMultisetOfNullableElement`, alongside 
`parquetRoundTripsNullMultiset` for a null multiset value in an optional column.
   
   **"What does this check?" on the old `testMultisetValueIsARequiredInt`** — 
you were right to ask, because it checked nothing. It asserted 
`isValueRequired()` on `SCHEMA`, a constant declared ten lines above in the 
same file, so it restated the test's own setup and would have passed with the 
production change reverted. It is replaced by 
`multisetConvertsToMapWithRequiredCount`, which puts a `MultisetType` through 
`FlinkSchemaUtil.convert` and asserts the resulting Iceberg `map<string, int>` 
has a required value — production behaviour, and the premise the writers depend 
on.
   
   **"Shall we add this to DynamicSink as well?"** — it needs a separate fix, 
and I would rather not grow this PR with it unless you prefer otherwise. 
`DataConverter#get` switches on `targetType.getTypeRoot()` and handles `ROW`, 
`ARRAY` and `MAP`; a `MULTISET` root reaches `default:` and throws 
`UnsupportedOperationException("Not a supported type: ...")` 
(`DataConverter.java:126`). So it is a different failure in a different place, 
reached from `TableUpdater` when a schema conversion is needed, and 
`CompareSchemasVisitor` carries a comment saying it must stay in sync with 
`DataConverter#get`, so the two want changing together. Happy to open a 
follow-up issue and PR for it, or to fold it in here if you would rather have 
one change.
   
   ### @Guosmilesmile
   
   **Pattern-matching `instanceof`** adopted — `FlinkMultisets#asMapType` uses 
`if (logicalType instanceof MultisetType multisetType)`. The module compiles at 
`options.release = 17`, so it is available.
   
   **`createTempDirectory` leaking a directory** — fixed, the ORC test uses 
`@TempDir` with the `File.createTempFile("junit", null, temp.toFile())` idiom 
the other tests in this package use. The Parquet and Avro tests need no file at 
all now, since `InMemoryOutputFile.toInputFile()` reads straight back.
   
   **Round-trip UT** — that is what all five write tests are now.
   
   **Whether `IntType` should be nullable, versus `AvroSchemaConverter` 
L606-609.** Good catch, and the honest answer is that both work, so this is 
about which one describes the type correctly rather than a bug:
   
   - `FlinkTypeToType#visit(MultisetType)` produces 
`Types.MapType.ofRequired(...)`, so the Iceberg schema the writers must match 
has a **required** count. `new IntType(false)` says the same thing, which is 
why I used it.
   - I checked whether it is load-bearing by flipping it to `new IntType()` and 
re-running: all six tests still pass. The writers use the logical type to 
choose a value writer, and a multiset count is never null, so nullability does 
not change what is written.
   
   So `AvroSchemaConverter`'s `new IntType()` is looser than the schema rather 
than wrong, and I have left it alone — it is a read-side schema conversion and 
out of scope here. If you would rather the two agreed, I am happy to tighten it 
in a follow-up.
   
   ### Verification
   
   `TestFlinkMultisetWrite` is 6 tests covering Parquet, Avro and ORC. **Five 
of the six fail without the production change**, with `Invalid map: 
MULTISET<STRING NOT NULL> is not a map` from the Parquet and Avro visitors and 
`ClassCastException: MultisetType cannot be cast to MapType` from ORC; the 
sixth is the conversion test, which passes either way because `FlinkTypeToType` 
already worked.
   
   Since the change now touches shared visitors used by every Flink writer, I 
ran the whole package: **711 tests in `org.apache.iceberg.flink.data`, 0 
failures**. `spotlessCheck`, `checkstyleMain` and `checkstyleTest` all pass. 
The three errorprone warnings that show on this module are on lines identical 
to `main` and are not from this change.
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to