Skip to content

Commit 8676144

Browse files
committed
address gemini comments
1 parent e74dd60 commit 8676144

2 files changed

Lines changed: 13 additions & 16 deletions

File tree

sdks/python/apache_beam/yaml/yaml_mapping.py

Lines changed: 12 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -20,8 +20,6 @@
2020
import datetime
2121
import itertools
2222
import re
23-
import threading
24-
import uuid
2523
from collections import abc
2624
from collections.abc import Callable
2725
from collections.abc import Collection
@@ -241,18 +239,13 @@ def process(self, element):
241239

242240
class JsMapToFieldsDoFn(beam.DoFn):
243241
def __init__(self, fields, original_fields, input_schema):
244-
self.fields = fields
245-
self.original_fields = original_fields
246-
self.input_schema = input_schema
247242
self.ctx = None
248243
self.field_funcs = {}
249244
self.passthrough_fields = []
250245

251-
def setup(self):
252-
self.ctx = MiniRacer()
253246
script = []
254-
for name, expr in self.fields.items():
255-
if isinstance(expr, str) and expr in self.input_schema:
247+
for name, expr in fields.items():
248+
if isinstance(expr, str) and expr in input_schema:
256249
self.passthrough_fields.append((name, expr))
257250
continue
258251

@@ -261,9 +254,9 @@ def setup(self):
261254

262255
if 'expression' in expr:
263256
e = expr['expression']
264-
code = f"var func_{name} = (__row__) => {{ " + " ".join([
265-
f"const {n} = __row__.{n};" for n in self.original_fields if n in e
266-
]) + f" return ({e}); }}"
257+
code = f"var func_{name} = (__row__) => {{ " + " ".join(
258+
[f"const {n} = __row__.{n};"
259+
for n in original_fields if n in e]) + f" return ({e}); }}"
267260
script.append(code)
268261
self.field_funcs[name] = f"func_{name}"
269262
elif 'callable' in expr:
@@ -277,8 +270,12 @@ def setup(self):
277270
script.append(udf_code)
278271
self.field_funcs[name] = func_name
279272

280-
if script:
281-
self.ctx.eval("\n".join(script))
273+
self.script = "\n".join(script) if script else None
274+
275+
def setup(self):
276+
self.ctx = MiniRacer()
277+
if self.script:
278+
self.ctx.eval(self.script)
282279

283280
def process(self, element):
284281
row_as_dict = py_value_to_js_dict(element)
@@ -308,7 +305,7 @@ def _get_javascript_udf_code(
308305

309306
if MiniRacer is None:
310307
raise ValueError(
311-
"JavaScript mapping functions require the 'mini-racer' package to be installed."
308+
"JavaScript mapping functions require the 'py-mini-racer' package to be installed."
312309
)
313310

314311
udf_code = None

sdks/python/setup.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -636,7 +636,7 @@ def get_portability_package_data():
636636
'docstring-parser>=0.15,<1.0',
637637
'jinja2>=3.0,<3.2',
638638
'virtualenv-clone>=0.5,<1.0',
639-
'mini-racer',
639+
'py-mini-racer',
640640
'jsonschema>=4.0.0,<5.0.0',
641641
] + dataframe_dependency,
642642
# Keep the following dependencies in line with what we test against

0 commit comments

Comments
 (0)