-
Notifications
You must be signed in to change notification settings - Fork 46
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
- Loading branch information
Showing
1 changed file
with
55 additions
and
0 deletions.
There are no files selected for viewing
55 changes: 55 additions & 0 deletions
55
centurion-avro/src/main/kotlin/com/sksamuel/centurion/avro/io/AvroWriter.kt
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,55 @@ | ||
package com.sksamuel.centurion.avro.io | ||
|
||
import org.apache.avro.Schema | ||
import org.apache.avro.generic.GenericDatumWriter | ||
import org.apache.avro.generic.GenericRecord | ||
import org.apache.avro.io.DatumWriter | ||
import org.apache.avro.io.EncoderFactory | ||
import java.io.OutputStream | ||
|
||
///** | ||
// * An [AvroWriter] will write [GenericRecord]s to an output stream. | ||
// * | ||
// * There are three implementations of this stream | ||
// * - a Data stream, | ||
// * - a Binary stream | ||
// * - a Json stream | ||
// * | ||
// * See the methods on the companion object to create instances of each | ||
// * of these types of stream. | ||
// */ | ||
|
||
/** | ||
* Creates an [AvroBinaryWriterFactory] for a given schema which can then be used | ||
* to create [AvroBinaryWriter]s. These writers share a thread safe [DatumWriter]. | ||
*/ | ||
class AvroBinaryWriterFactory(schema: Schema, private val factory: EncoderFactory) { | ||
constructor(schema: Schema) : this(schema, EncoderFactory.get()) | ||
|
||
private val datumWriter = GenericDatumWriter<GenericRecord>(schema) | ||
fun writer(output: OutputStream): AvroBinaryWriter { | ||
return AvroBinaryWriter(datumWriter, output, factory) | ||
} | ||
} | ||
|
||
class AvroBinaryWriter( | ||
private val datumWriter: DatumWriter<GenericRecord>, | ||
private val output: OutputStream, | ||
private val factory: EncoderFactory, | ||
) : AutoCloseable { | ||
|
||
private val encoder = factory.binaryEncoder(output, null) | ||
|
||
fun write(record: GenericRecord) { | ||
datumWriter.write(record, encoder) | ||
} | ||
|
||
fun flush() { | ||
encoder.flush() | ||
} | ||
|
||
override fun close() { | ||
flush() | ||
output.close() | ||
} | ||
} |