@@ -15,7 +15,7 @@ import io.airbyte.cdk.load.message.DestinationFile
15
15
import io.airbyte.cdk.load.message.DestinationFileStreamComplete
16
16
import io.airbyte.cdk.load.message.DestinationFileStreamIncomplete
17
17
import io.airbyte.cdk.load.message.DestinationRecord
18
- import io.airbyte.cdk.load.message.DestinationRecordAirbyteValue
18
+ import io.airbyte.cdk.load.message.DestinationRecordRaw
19
19
import io.airbyte.cdk.load.message.DestinationRecordStreamComplete
20
20
import io.airbyte.cdk.load.message.DestinationRecordStreamIncomplete
21
21
import io.airbyte.cdk.load.message.DestinationStreamAffinedMessage
@@ -80,7 +80,7 @@ class DefaultInputConsumerTask(
80
80
// Required by new interface
81
81
@Named(" recordQueue" )
82
82
private val recordQueueForPipeline :
83
- PartitionedQueue <Reserved <PipelineEvent <StreamKey , DestinationRecordAirbyteValue >>>,
83
+ PartitionedQueue <Reserved <PipelineEvent <StreamKey , DestinationRecordRaw >>>,
84
84
private val loadPipeline : LoadPipeline ? = null ,
85
85
private val partitioner : InputPartitioner ,
86
86
private val openStreamQueue : QueueWriter <DestinationStream >
@@ -158,7 +158,7 @@ class DefaultInputConsumerTask(
158
158
val manager = syncManager.getStreamManager(stream)
159
159
when (val message = reserved.value) {
160
160
is DestinationRecord -> {
161
- val record = message.asRecordMarshaledToAirbyteValue ()
161
+ val record = message.asDestinationRecordRaw ()
162
162
manager.incrementReadCount()
163
163
val pipelineMessage =
164
164
PipelineMessage (
@@ -310,7 +310,7 @@ interface InputConsumerTaskFactory {
310
310
311
311
// Required by new interface
312
312
recordQueueForPipeline :
313
- PartitionedQueue <Reserved <PipelineEvent <StreamKey , DestinationRecordAirbyteValue >>>,
313
+ PartitionedQueue <Reserved <PipelineEvent <StreamKey , DestinationRecordRaw >>>,
314
314
loadPipeline : LoadPipeline ? ,
315
315
partitioner : InputPartitioner ,
316
316
openStreamQueue : QueueWriter <DestinationStream >,
@@ -333,7 +333,7 @@ class DefaultInputConsumerTaskFactory(
333
333
334
334
// Required by new interface
335
335
recordQueueForPipeline :
336
- PartitionedQueue <Reserved <PipelineEvent <StreamKey , DestinationRecordAirbyteValue >>>,
336
+ PartitionedQueue <Reserved <PipelineEvent <StreamKey , DestinationRecordRaw >>>,
337
337
loadPipeline : LoadPipeline ? ,
338
338
partitioner : InputPartitioner ,
339
339
openStreamQueue : QueueWriter <DestinationStream >,
0 commit comments