File tree Expand file tree Collapse file tree 1 file changed +2
-2
lines changed
kcbq-connector/src/main/java/com/wepay/kafka/connect/bigquery Expand file tree Collapse file tree 1 file changed +2
-2
lines changed Original file line number Diff line number Diff line change @@ -154,10 +154,10 @@ private PartitionedTableId getRecordTable(SinkRecord record) {
154
154
private RowToInsert getRecordRow (SinkRecord record ) {
155
155
Map <String , Object > convertedRecord = recordConverter .convertRecord (record .valueSchema (), record .value ());
156
156
if (config .getBoolean (config .INCLUDE_KAFKA_KEY_CONFIG )) {
157
- convertedRecord .put (config .KAFKA_KEY_FIELD_NAME_CONFIG , recordConverter .convertRecord (record .keySchema (), record .key ()));
157
+ convertedRecord .put (config .getString ( config . KAFKA_KEY_FIELD_NAME_CONFIG ) , recordConverter .convertRecord (record .keySchema (), record .key ()));
158
158
}
159
159
if (config .getBoolean (config .INCLUDE_KAFKA_DATA_CONFIG )) {
160
- convertedRecord .put (config .KAFKA_DATA_FIELD_NAME_CONFIG , KafkaDataBuilder .getKafkaDataRecord (record ));
160
+ convertedRecord .put (config .getString ( config . KAFKA_DATA_FIELD_NAME_CONFIG ) , KafkaDataBuilder .getKafkaDataRecord (record ));
161
161
}
162
162
if (config .getBoolean (config .SANITIZE_FIELD_NAME_CONFIG )) {
163
163
convertedRecord = FieldNameSanitizer .replaceInvalidKeys (convertedRecord );
You can’t perform that action at this time.
0 commit comments