Hi all, I am writing a custom source (Iceberg in t...
# ingestion
m
Hi all, I am writing a custom source (Iceberg in this case. I know it's on the roadmap, but I need it now and I'm using this to understand the datahub internals) and I am having problems adding a
MetadataChangeEvent
with a
SchemaMetadata
aspect. It looks like something is rejected by the Avro validator, but it doesn't tell me what. Is there a trick to figure what exactly is incompatible with the schema?
Copy code
File "/datahub/metadata-ingestion/src/datahub/cli/ingest_cli.py", line 82, in run
    pipeline.run()
File "/datahub/metadata-ingestion/src/datahub/ingestion/run/pipeline.py", line 157, in run
    for record_envelope in self.transform(record_envelopes):
File "/datahub/metadata-ingestion/src/datahub/ingestion/extractor/mce_extractor.py", line 46, in get_records
    raise ValueError(

ValueError: source produced an invalid metadata work unit: MetadataChangeEventClass(...
h
Hi @modern-monitor-81461, you typically run into this when the mce is missing some required members or if there is an issue with the values it contains. One way to troubleshoot this is (1) call validate method on the mce that you have crated (my_mce.validate()) and (2) step through the debugger and see where the method is returning false from.
m
@helpful-optician-78938 Ok, that was a good suggestion and I did that. I now know what is wrong, but I don't understand why and it probably has to do with my lack of DataHub knowledge:
Copy code
field = SchemaField(
            fieldPath="test",
            nativeDataType="bool",
            type=BooleanTypeClass,
        )
        print("BEFORE FIELD VALIDATE")
        field.validate()
        print("AFTER FIELD VALIDATE")
I create a
SchemaField
and I set its type to be a
boolean
. When I validate, the
type
attribute is invalid and I don't understand why. I have modified the avrojson.py code like this:
Copy code
elif schema_type in ['record', 'error', 'request']:
            #####  start of test
            if (isinstance(datum, dict) or isinstance(datum, DictWrapper)):
                result = []
                for f in expected_schema.fields:
                    print(f"Validating {f.name} is {self.validate(f.type, datum.get(f.name), skip_logical_types)}")
                    testResult = self.validate(f.type, datum.get(f.name), skip_logical_types)
                    result.append(testResult)
                return False not in result
            ##### end of test

            # return ((isinstance(datum, dict) or isinstance(datum, DictWrapper)) and
            #         False not in
            #         [self.validate(f.type, datum.get(f.name), skip_logical_types) for f in expected_schema.fields])
and I get the following output:
Copy code
BEFORE FIELD VALIDATE
Validating fieldPath is True
Validating jsonPath is True
Validating nullable is True
Validating description is True
Validating type is False
Validating nativeDataType is True
Validating recursive is True
Validating globalTags is True
Validating glossaryTerms is True
Validating isPartOfKey is True
Validating jsonProps is True
AFTER FIELD VALIDATE
All in all, what I'm trying to do is read an Apache Iceberg schema, get all the columns and create a
SchemaField
for each one.
fieldPath
would simply be the name of the column,
nativeDataType
would be the native Iceberg type and
type
would be what I map the column to in DataHub. Is my logic good? If so, why is my code not passing validation?
h
Hi @modern-monitor-81461, You should actually use an instance of BooleanTypeClass for the type, not the plain class name, wrapped inside the
SchemaFieldDataTypeClass
. Can you try changing it to
type=SchemaFieldDataTypeClass(type=BooleanTypeClass())
? Here is an example.
m
That did the trick, thank you sir! 😃
👍 1