Skip to content
Closed
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
65 changes: 31 additions & 34 deletions Sources/SparkConnect/ArrowReader.swift
Original file line number Diff line number Diff line change
Expand Up @@ -274,23 +274,22 @@ public class ArrowReader { // swiftlint:disable:this type_body_length
let message = org_apache_arrow_flatbuf_Message.getRootAsMessage(bb: dataBuffer)
switch message.headerType {
case .recordbatch:
do {
let rbMessage = message.header(type: org_apache_arrow_flatbuf_RecordBatch.self)!
let recordBatch = try loadRecordBatch(
rbMessage,
schema: schemaMessage!,
arrowSchema: result.schema!,
data: input,
messageEndOffset: (Int64(offset) + Int64(length))
).get()
let rbMessage = message.header(type: org_apache_arrow_flatbuf_RecordBatch.self)!
let recordBatchResult = try loadRecordBatch(
rbMessage,
schema: schemaMessage!,
arrowSchema: result.schema!,
data: input,
messageEndOffset: (Int64(offset) + Int64(length))
)
switch recordBatchResult {
case .success(let recordBatch):
result.batches.append(recordBatch)
offset += Int(message.bodyLength + Int64(length))
length = getUInt32(input, offset: offset)
} catch let error as ArrowError {
case .failure(let error):
return .failure(error)
} catch {
return .failure(.unknownError("Unexpected error: \(error)"))
}
offset += Int(message.bodyLength + Int64(length))
length = getUInt32(input, offset: offset)
case .schema:
schemaMessage = message.header(type: org_apache_arrow_flatbuf_Schema.self)!
let schemaResult = loadSchema(schemaMessage!)
Expand Down Expand Up @@ -363,20 +362,19 @@ public class ArrowReader { // swiftlint:disable:this type_body_length
let message = org_apache_arrow_flatbuf_Message.getRootAsMessage(bb: mbb)
switch message.headerType {
case .recordbatch:
do {
let rbMessage = message.header(type: org_apache_arrow_flatbuf_RecordBatch.self)!
let recordBatch = try loadRecordBatch(
rbMessage,
schema: footer.schema!,
arrowSchema: result.schema!,
data: fileData,
messageEndOffset: messageEndOffset
).get()
let rbMessage = message.header(type: org_apache_arrow_flatbuf_RecordBatch.self)!
let recordBatchResult = try loadRecordBatch(
rbMessage,
schema: footer.schema!,
arrowSchema: result.schema!,
data: fileData,
messageEndOffset: messageEndOffset
)
switch recordBatchResult {
case .success(let recordBatch):
result.batches.append(recordBatch)
} catch let error as ArrowError {
case .failure(let error):
return .failure(error)
} catch {
return .failure(.unknownError("Unexpected error: \(error)"))
}
default:
return .failure(.unknownError("Unhandled header type: \(message.headerType)"))
Expand Down Expand Up @@ -429,17 +427,16 @@ public class ArrowReader { // swiftlint:disable:this type_body_length
}
case .recordbatch:
let rbMessage = message.header(type: org_apache_arrow_flatbuf_RecordBatch.self)!
do {
let recordBatch = try loadRecordBatch(
rbMessage, schema: result.messageSchema!, arrowSchema: result.schema!,
data: dataBody, messageEndOffset: 0
).get()
let recordBatchResult = try loadRecordBatch(
rbMessage, schema: result.messageSchema!, arrowSchema: result.schema!,
data: dataBody, messageEndOffset: 0
)
switch recordBatchResult {
case .success(let recordBatch):
result.batches.append(recordBatch)
return .success(())
} catch let error as ArrowError {
case .failure(let error):
return .failure(error)
} catch {
return .failure(.unknownError("Unexpected error: \(error)"))
}
default:
return .failure(.unknownError("Unhandled header type: \(message.headerType)"))
Expand Down
Loading