Skip to content

Commit 3dcf794

Browse files
committed
feat(dataframe): add withColumnRenamed method
1 parent b8673a9 commit 3dcf794

3 files changed

Lines changed: 79 additions & 0 deletions

File tree

native/src/lib.rs

Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -308,6 +308,26 @@ pub extern "system" fn Java_org_apache_datafusion_DataFrame_dropColumns<'local>(
308308
})
309309
}
310310

311+
#[no_mangle]
312+
pub extern "system" fn Java_org_apache_datafusion_DataFrame_renameColumn<'local>(
313+
mut env: JNIEnv<'local>,
314+
_class: JClass<'local>,
315+
handle: jlong,
316+
old_name: JString<'local>,
317+
new_name: JString<'local>,
318+
) -> jlong {
319+
try_unwrap_or_throw(&mut env, 0, |env| -> JniResult<jlong> {
320+
if handle == 0 {
321+
return Err("DataFrame handle is null".into());
322+
}
323+
let df = unsafe { &*(handle as *const DataFrame) }.clone();
324+
let old: String = env.get_string(&old_name)?.into();
325+
let new: String = env.get_string(&new_name)?.into();
326+
let new_df = df.with_column_renamed(&old, &new)?;
327+
Ok(Box::into_raw(Box::new(new_df)) as jlong)
328+
})
329+
}
330+
311331
#[no_mangle]
312332
pub extern "system" fn Java_org_apache_datafusion_DataFrame_writeParquetWithOptions<'local>(
313333
mut env: JNIEnv<'local>,

src/main/java/org/apache/datafusion/DataFrame.java

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -163,6 +163,14 @@ public DataFrame dropColumns(String... columnNames) {
163163
return new DataFrame(dropColumns(nativeHandle, columnNames));
164164
}
165165

166+
/** Rename a column. The receiver remains usable and must still be closed independently. */
167+
public DataFrame withColumnRenamed(String oldName, String newName) {
168+
if (nativeHandle == 0) {
169+
throw new IllegalStateException("DataFrame is closed or already collected");
170+
}
171+
return new DataFrame(renameColumn(nativeHandle, oldName, newName));
172+
}
173+
166174
/**
167175
* Materialize this DataFrame as Parquet at {@code path}. The path is treated as a directory
168176
* unless overridden via {@link ParquetWriteOptions#singleFileOutput(boolean)}. The receiver
@@ -221,6 +229,8 @@ public void close() {
221229

222230
private static native long dropColumns(long handle, String[] columnNames);
223231

232+
private static native long renameColumn(long handle, String oldName, String newName);
233+
224234
private static native void writeParquetWithOptions(
225235
long handle,
226236
String path,

src/test/java/org/apache/datafusion/DataFrameTransformationsTest.java

Lines changed: 49 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -294,4 +294,53 @@ void dropColumnsSilentlyIgnoresUnknownNames() throws Exception {
294294
root.getSchema().getFields().stream().map(f -> f.getName()).toArray(String[]::new));
295295
}
296296
}
297+
298+
@Test
299+
void withColumnRenamedChangesColumnName() throws Exception {
300+
try (BufferAllocator allocator = new RootAllocator();
301+
SessionContext ctx = new SessionContext();
302+
DataFrame source = ctx.sql("SELECT 1 AS a, 2 AS b");
303+
DataFrame renamed = source.withColumnRenamed("a", "alpha");
304+
ArrowReader reader = renamed.collect(allocator)) {
305+
assertTrue(reader.loadNextBatch());
306+
VectorSchemaRoot root = reader.getVectorSchemaRoot();
307+
assertArrayEquals(
308+
new String[] {"alpha", "b"},
309+
root.getSchema().getFields().stream().map(f -> f.getName()).toArray(String[]::new));
310+
}
311+
}
312+
313+
@Test
314+
void withColumnRenamedIsNonDestructive() throws Exception {
315+
try (BufferAllocator allocator = new RootAllocator();
316+
SessionContext ctx = new SessionContext();
317+
DataFrame source = ctx.sql("SELECT 1 AS a, 2 AS b")) {
318+
try (DataFrame renamed = source.withColumnRenamed("a", "alpha")) {
319+
assertEquals(1L, renamed.count());
320+
}
321+
try (DataFrame again = source.select("a");
322+
ArrowReader reader = again.collect(allocator)) {
323+
assertTrue(reader.loadNextBatch());
324+
VectorSchemaRoot root = reader.getVectorSchemaRoot();
325+
assertArrayEquals(
326+
new String[] {"a"},
327+
root.getSchema().getFields().stream().map(f -> f.getName()).toArray(String[]::new));
328+
}
329+
}
330+
}
331+
332+
@Test
333+
void withColumnRenamedUnknownColumnIsNoOp() throws Exception {
334+
try (BufferAllocator allocator = new RootAllocator();
335+
SessionContext ctx = new SessionContext();
336+
DataFrame df = ctx.sql("SELECT 1 AS x");
337+
DataFrame renamed = df.withColumnRenamed("not_a_column", "y");
338+
ArrowReader reader = renamed.collect(allocator)) {
339+
assertTrue(reader.loadNextBatch());
340+
VectorSchemaRoot root = reader.getVectorSchemaRoot();
341+
assertArrayEquals(
342+
new String[] {"x"},
343+
root.getSchema().getFields().stream().map(f -> f.getName()).toArray(String[]::new));
344+
}
345+
}
297346
}

0 commit comments

Comments
 (0)