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]