davidradl commented on code in PR #56:
URL:
https://github.com/apache/flink-connector-http/pull/56#discussion_r4218934203
##########
flink-connector-http/src/main/java/org/apache/flink/connector/http/table/lookup/HttpTableLookupFunction.java:
##########
@@ -107,90 +103,72 @@ public Collection<RowData> lookup(RowData keyRow) {
HttpRowDataWrapper httpRowDataWrapper = client.pull(keyRow);
Collection<RowData> httpCollector = httpRowDataWrapper.getData();
- int physicalArity = -1;
-
- GenericRowData producedRow = null;
- // grab the actual data if there is any from the response and populate
the producedRow with
- // it
- if (!httpCollector.isEmpty()) {
- GenericRowData physicalRow = (GenericRowData)
httpCollector.iterator().next();
- physicalArity = physicalRow.getArity();
- producedRow =
- new GenericRowData(physicalRow.getRowKind(), physicalArity
+ metadataArity);
- // Build a map of lookup table field names to their typed values
from keyRow
- // Only process top-level single-value join keys
- Map<String, Object> joinKeyValues = new HashMap<>();
- for (LookupSchemaEntry<RowData> entry :
lookupRow.getLookupEntries()) {
- // Only handle top-level single value entries (not nested
RowTypeLookupSchemaEntry)
- if (entry instanceof RowDataSingleValueLookupSchemaEntry) {
- RowDataSingleValueLookupSchemaEntry singleEntry =
- (RowDataSingleValueLookupSchemaEntry) entry;
- try {
- // Get the typed value directly from keyRow using
fieldGetter
- Object typedValue =
singleEntry.fieldGetter.getFieldOrNull(keyRow);
- if (typedValue != null) {
- // Get the lookup table field name from LookupArg
- List<LookupArg> lookupArgs =
entry.convertToLookupArg(keyRow);
- for (LookupArg lookupArg : lookupArgs) {
- // Map lookup table field name to typed value
- joinKeyValues.put(lookupArg.getArgName(),
typedValue);
- }
+ if (httpCollector.isEmpty() && metadataArity == 0) {
+ // A response without data is a lookup miss unless the row carries
metadata columns.
+ return Collections.emptyList();
+ }
+ final GenericRowData physicalRow =
+ httpCollector.isEmpty()
+ ? new GenericRowData(
+ RowKind.INSERT,
physicalRowDataType.getChildren().size())
+ : (GenericRowData) httpCollector.iterator().next();
+ final int physicalArity = physicalRow.getArity();
+ GenericRowData producedRow =
+ new GenericRowData(physicalRow.getRowKind(), physicalArity +
metadataArity);
+ // Build a map of lookup table field names to their typed values from
keyRow
+ // Only process top-level single-value join keys
+ Map<String, Object> joinKeyValues = new HashMap<>();
+ for (LookupSchemaEntry<RowData> entry : lookupRow.getLookupEntries()) {
+ // Only handle top-level single value entries (not nested
RowTypeLookupSchemaEntry)
+ if (entry instanceof RowDataSingleValueLookupSchemaEntry) {
+ RowDataSingleValueLookupSchemaEntry singleEntry =
+ (RowDataSingleValueLookupSchemaEntry) entry;
+ try {
+ // Get the typed value directly from keyRow using
fieldGetter
+ Object typedValue =
singleEntry.fieldGetter.getFieldOrNull(keyRow);
+ if (typedValue != null) {
+ // Get the lookup table field name from LookupArg
+ List<LookupArg> lookupArgs =
entry.convertToLookupArg(keyRow);
+ for (LookupArg lookupArg : lookupArgs) {
+ // Map lookup table field name to typed value
+ joinKeyValues.put(lookupArg.getArgName(),
typedValue);
}
- } catch (Exception e) {
- log.warn(
- "Failed to extract join key value for field:
{}",
- entry.getFieldName(),
- e);
}
+ } catch (Exception e) {
+ log.warn(
+ "Failed to extract join key value for field: {}",
+ entry.getFieldName(),
+ e);
}
}
+ }
- // Get physical row field names to match positions
- List<String> physicalFieldNames =
-
TableSourceHelper.getFieldNames(physicalRowDataType.getLogicalType());
-
- // Copy fields from physicalRow to producedRow, populating null
join keys
- for (int pos = 0; pos < physicalArity; pos++) {
- Object value = physicalRow.getField(pos);
- String fieldName = physicalFieldNames.get(pos);
- // If field is null and it's a join key, populate from keyRow
- if (value == null && !joinKeyValues.isEmpty() && pos <
physicalFieldNames.size()) {
- if (joinKeyValues.containsKey(fieldName)) {
- value = joinKeyValues.get(fieldName);
- if (log.isDebugEnabled()) {
- log.debug(
- "Lookup processing found a value null for
join key {}, replacing value with the {}",
- fieldName,
- value.toString());
- }
+ // Get physical row field names to match positions
+ List<String> physicalFieldNames =
+
TableSourceHelper.getFieldNames(physicalRowDataType.getLogicalType());
+
+ // Copy fields from physicalRow to producedRow, populating null join
keys
+ for (int pos = 0; pos < physicalArity; pos++) {
+ Object value = physicalRow.getField(pos);
+ String fieldName = physicalFieldNames.get(pos);
+ // If field is null and it's a join key, populate from keyRow
+ if (value == null && !joinKeyValues.isEmpty() && pos <
physicalFieldNames.size()) {
+ if (joinKeyValues.containsKey(fieldName)) {
+ value = joinKeyValues.get(fieldName);
Review Comment:
I am not convinced with should override the id column with the probe column
value if we get a null back.
So the table to be enriched has id column an event with id=1 comes in, we
issue a lookup using id=1, but the response gives us a null value in id, since
the id here is a response field. In this case I think we should honour the
null. In the same way we would honour values that were no null and did not
match the request join key value.
We could add to the docs to say something like: if the user wants to process
field names that are the same and the same type and exist in the the API
request and response, it is possible to define one column; in this case the
response value will be populated after the API call. If the user would like to
keep the request and response values separate then the response column should
be defined with the name in the response, the path, query param or body field
should be mapped from a differently named column name. In this way you can keep
both the request and the response values in SQL columns. This is the pattern we
use.
--
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]